Skip to content

Commit 659efb6

Browse files
committed
Experimental python3 support
1 parent 47cff11 commit 659efb6

37 files changed

Lines changed: 244 additions & 139 deletions

examples/simple_crawler/crawler.py

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -86,16 +86,16 @@ def run(self, params):
8686

8787
collection = connections.mongodb_jobs.simple_crawler_urls
8888

89-
print
90-
print "Crawl stats"
91-
print "==========="
92-
print "URLs queued: %s" % collection.find().count()
93-
print "URLs successfully crawled: %s" % collection.find({"fetched_date": {"$exists": True}}).count()
94-
print "URLs redirected: %s" % collection.find({"redirected_to": {"$exists": True}}).count()
95-
print "Bytes fetched: %s" % (list(collection.aggregate(
89+
print()
90+
print( "Crawl stats")
91+
print( "===========")
92+
print( "URLs queued: %s" % collection.find().count())
93+
print( "URLs successfully crawled: %s" % collection.find({"fetched_date": {"$exists": True}}).count())
94+
print( "URLs redirected: %s" % collection.find({"redirected_to": {"$exists": True}}).count())
95+
print( "Bytes fetched: %s" % (list(collection.aggregate(
9696
{"$group": {"_id": None, "sum": {"$sum": "$html_length"}}}
97-
)) or [{}])[0].get("sum", 0)
98-
print
97+
)) or [{}])[0].get("sum", 0))
98+
print()
9999

100100

101101
class Reset(Task):

mrq/basetasks/utils.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
from __future__ import print_function
12
from mrq.task import Task
23
from mrq.queue import Queue
34
from bson import ObjectId
@@ -50,7 +51,7 @@ def build_query(self):
5051
if self.params.get("params"):
5152
params_dict = json.loads(self.params.get("params")) # pylint: disable=no-member
5253

53-
for key in params_dict.keys():
54+
for key in list(params_dict.keys()):
5455
query["params.%s" % key] = params_dict[key]
5556

5657
return query
@@ -76,7 +77,7 @@ def perform_action(self, action, query, destination_queue):
7677
else:
7778

7879
tasks_defs = get_current_config().get("tasks", {})
79-
tasks_ttls = [cfg.get("result_ttl", 0) for cfg in tasks_defs.values()]
80+
tasks_ttls = [cfg.get("result_ttl", 0) for cfg in list(tasks_defs.values())]
8081

8182
result_ttl = max([default_job_timeout] + tasks_ttls)
8283

@@ -92,7 +93,7 @@ def perform_action(self, action, query, destination_queue):
9293
# In this case we could also loose some jobs that were queued after
9394
# the MongoDB update. They will be "lost" and requeued later like the other case
9495
# after the Redis BLPOP
95-
if query.keys() == ["queue"]:
96+
if list(query.keys()) == ["queue"]:
9697
Queue(query["queue"]).empty()
9798

9899
elif action in ("requeue", "requeue_retry"):
@@ -135,6 +136,6 @@ def perform_action(self, action, query, destination_queue):
135136
Queue(destination_queue or queue).enqueue_job_ids(
136137
[str(x) for x in jobs_by_queue[queue]])
137138

138-
print stats
139+
print(stats)
139140

140141
return stats

mrq/bin/mrq_run.py

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
11
#!/usr/bin/env python
2+
from __future__ import print_function
3+
24
import os
35

46
# Needed to make getaddrinfo() work in pymongo on Mac OS X
@@ -40,13 +42,13 @@ def main():
4042
# mrq-run taskpath a 1 b 2 => {"a": "1", "b": "2"}
4143
for group in utils.group_iter(cfg["taskargs"], n=2):
4244
if len(group) != 2:
43-
print "Number of arguments wasn't even"
45+
print("Number of arguments wasn't even")
4446
sys.exit(1)
4547
params[group[0]] = group[1]
4648

4749
if cfg["queue"]:
4850
ret = queue_job(cfg["taskpath"], params, queue=cfg["queue"])
49-
print ret
51+
print(ret)
5052
else:
5153
worker_class = load_class_by_path(cfg["worker_class"])
5254
job = worker_class.job_class(None)
@@ -58,7 +60,7 @@ def main():
5860
job.datestarted = datetime.datetime.utcnow()
5961
set_current_job(job)
6062
ret = job.perform()
61-
print json_stdlib.dumps(ret, cls=MongoJSONEncoder) # pylint: disable=no-member
63+
print(json_stdlib.dumps(ret, cls=MongoJSONEncoder)) # pylint: disable=no-member
6264

6365
# This shouldn't be needed as the process will exit and close any remaining sockets
6466
# connections.redis.connection_pool.disconnect()

mrq/bin/mrq_worker.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,15 +8,19 @@
88
os.environ["GEVENT_RESOLVER"] = "ares"
99

1010
from gevent import monkey
11-
monkey.patch_all()
11+
monkey.patch_all(subprocess=False)
1212

1313
import sys
1414
import tempfile
1515
import signal
16-
import subprocess32 as subprocess
1716
import psutil
1817
import argparse
1918

19+
try:
20+
import subprocess32 as subprocess
21+
except:
22+
import subprocess
23+
2024
sys.path.insert(0, os.getcwd())
2125

2226
from mrq import config

mrq/config.py

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
from __future__ import print_function
12
import argparse
23
import os
34
import sys
@@ -325,7 +326,7 @@ def add_parser_args(parser, config_type):
325326
'--report_file',
326327
default="",
327328
action='store',
328-
type=unicode,
329+
type=str,
329330
help='Filepath of a json dump of the worker status. Disabled if none')
330331

331332
parser.add_argument(
@@ -402,7 +403,7 @@ def get_config(
402403
# line
403404
from_args = {}
404405
if "args" in sources:
405-
for k, v in parser.parse_args().__dict__.iteritems():
406+
for k, v in parser.parse_args().__dict__.items():
406407
if default_config[k] != v:
407408
from_args[k] = v
408409

@@ -422,7 +423,7 @@ def get_config(
422423
sys.path.insert(0, os.path.dirname(config_file))
423424
config_module = __import__(os.path.basename(config_file.replace(".py", "")))
424425
sys.path.pop(0)
425-
for k, v in config_module.__dict__.iteritems():
426+
for k, v in config_module.__dict__.items():
426427

427428
# We only keep variables starting with an uppercase character.
428429
if k[0].isupper():
@@ -462,11 +463,11 @@ def print_profiling():
462463
atexit.register(print_profiling)
463464

464465
if merged_config["version"]:
465-
print "MRQ version: %s" % VERSION
466-
print "Python version: %s" % sys.version
466+
print("MRQ version: %s" % VERSION)
467+
print("Python version: %s" % sys.version)
467468
sys.exit(1)
468469

469470
if "no_import_patch" in from_args:
470-
print "WARNING: --no_import_patch will be deprecated in MRQ 1.0!"
471+
print("WARNING: --no_import_patch will be deprecated in MRQ 1.0!")
471472

472473
return merged_config

mrq/context.py

Lines changed: 13 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
1+
from future import standard_library
2+
standard_library.install_aliases()
3+
from builtins import next
4+
from builtins import map
15
from .logger import Logger
26
import gevent
37
import gevent.pool
4-
import urlparse
8+
import urllib.parse
59
import time
610
import pymongo
711
import traceback
@@ -108,12 +112,12 @@ def versiontuple(v):
108112
return tuple(map(int, (v.split("."))))
109113

110114
if attr.startswith("redis"):
111-
if type(config_obj) in [str, unicode]:
115+
if type(config_obj) in [str]:
112116

113117
import redis as pyredis
114118

115-
urlparse.uses_netloc.append('redis')
116-
redis_url = urlparse.urlparse(config_obj)
119+
urllib.parse.uses_netloc.append('redis')
120+
redis_url = urllib.parse.urlparse(config_obj)
117121

118122
log.info("%s: Connecting to Redis at %s..." %
119123
(attr, redis_url.hostname))
@@ -124,7 +128,8 @@ def versiontuple(v):
124128
db=int((redis_url.path or "").replace("/", "") or "0"),
125129
password=redis_url.password,
126130
max_connections=int(config.get("redis_max_connections")),
127-
timeout=int(config.get("redis_timeout"))
131+
timeout=int(config.get("redis_timeout")),
132+
decode_responses=True
128133
)
129134
return pyredis.StrictRedis(connection_pool=redis_pool)
130135

@@ -134,7 +139,7 @@ def versiontuple(v):
134139

135140
elif attr.startswith("mongodb"):
136141

137-
if type(config_obj) in [str, unicode]:
142+
if type(config_obj) in [str]:
138143

139144
if attr == "mongodb_logs" and config_obj == "1":
140145
return connections.mongodb_jobs
@@ -218,7 +223,7 @@ def inner_func(*args):
218223

219224
try:
220225
ret = func(*args)
221-
except Exception, exc:
226+
except Exception as exc:
222227
trace = traceback.format_exc()
223228
log.error("Error in subpool: %s \n%s" % (exc, trace))
224229
raise
@@ -271,7 +276,7 @@ def inner_func(*args):
271276

272277
try:
273278
ret = func(*args)
274-
except Exception, exc:
279+
except Exception as exc:
275280
trace = traceback.format_exc()
276281
log.error("Error in subpool: %s \n%s" % (exc, trace))
277282
raise

mrq/dashboard/app.py

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,6 @@
1+
from __future__ import print_function
2+
from future import standard_library
3+
standard_library.install_aliases()
14
from gevent import monkey
25
monkey.patch_all()
36

@@ -41,7 +44,7 @@
4144
@requires_auth
4245
def root():
4346
return render_template("index.html", MRQ_CONFIG={
44-
k: v for k, v in cfg.items() if k in WHITELISTED_MRQ_CONFIG_KEYS
47+
k: v for k, v in list(cfg.items()) if k in WHITELISTED_MRQ_CONFIG_KEYS
4548
})
4649

4750

@@ -149,10 +152,10 @@ def build_api_datatables_query(req):
149152
try:
150153
params_dict = json.loads(req.args.get("params"))
151154

152-
for key in params_dict.keys():
155+
for key in list(params_dict.keys()):
153156
query["params.%s" % key] = params_dict[key]
154157
except Exception as e: # pylint: disable=broad-except
155-
print "Error will converting form JSON: %s" % e
158+
print("Error will converting form JSON: %s" % e)
156159

157160
return query
158161

@@ -301,7 +304,7 @@ def api_job_traceback(job_id):
301304
@app.route('/api/jobaction', methods=["POST"])
302305
@requires_auth
303306
def api_job_action():
304-
params = {k: v for k, v in request.form.iteritems()}
307+
params = {k: v for k, v in request.form.items()}
305308
if params.get("status") and "-" in params.get("status"):
306309
params["status"] = params.get("status").split("-")
307310
return jsonify({"job_id": queue_job("mrq.basetasks.utils.JobAction",

mrq/job.py

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,7 @@
1+
from future import standard_library
2+
standard_library.install_aliases()
3+
from builtins import str
4+
from builtins import object
15
import datetime
26
from bson import ObjectId
37
import time
@@ -10,12 +14,12 @@
1014
from collections import defaultdict
1115
import traceback
1216
import sys
13-
import urlparse
17+
import urllib.parse
1418
import re
1519
import linecache
1620
import fnmatch
1721
import encodings
18-
import copy_reg
22+
import copyreg
1923
from . import context
2024

2125

@@ -54,7 +58,10 @@ def __init__(self, job_id, queue=None, start=False, fetch=False):
5458
if job_id is None:
5559
self.id = None
5660
else:
57-
self.id = ObjectId(job_id)
61+
if isinstance(job_id, bytes):
62+
self.id = ObjectId(job_id.decode('utf-8'))
63+
else:
64+
self.id = ObjectId(job_id)
5865

5966
self.data = None
6067
self.saved = True
@@ -463,10 +470,10 @@ def set_current_io(self, io_data):
463470
def trace_memory_clean_caches(self):
464471
""" Avoid polluting results with some builtin python caches """
465472

466-
urlparse.clear_cache()
473+
urllib.parse.clear_cache()
467474
re.purge()
468475
linecache.clearcache()
469-
copy_reg.clear_extension_cache()
476+
copyreg.clear_extension_cache()
470477

471478
if hasattr(fnmatch, "purge"):
472479
fnmatch.purge() # pylint: disable=no-member

mrq/logger.py

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,27 @@
1+
from __future__ import print_function
2+
from builtins import object
13

24
from collections import defaultdict
35
import datetime
4-
6+
import sys
7+
PY3 = sys.version_info > (3,)
58

69
def _encode_if_unicode(string):
10+
11+
if PY3:
12+
return string
13+
714
if isinstance(string, unicode):
815
return string.encode("utf-8", "replace")
916
else:
1017
return string
1118

1219

1320
def _decode_if_str(string):
21+
22+
if PY3:
23+
return str(string)
24+
1425
if isinstance(string, str):
1526
return string.decode("utf-8", "replace")
1627
else:
@@ -62,9 +73,9 @@ def log(self, level, *args, **kwargs):
6273

6374
if not self.quiet:
6475
try:
65-
print _encode_if_unicode(formatted)
76+
print(_encode_if_unicode(formatted))
6677
except UnicodeDecodeError:
67-
print formatted
78+
print(formatted)
6879

6980
if self.collection is False:
7081
return
@@ -88,10 +99,10 @@ def flush(self, w=0):
8899
inserts = [{
89100
"worker": k,
90101
"logs": "\n".join(v) + "\n"
91-
} for k, v in self.buffer["workers"].iteritems()] + [{
102+
} for k, v in self.buffer["workers"].items()] + [{
92103
"job": k,
93104
"logs": "\n".join(v) + "\n"
94-
} for k, v in self.buffer["jobs"].iteritems()]
105+
} for k, v in self.buffer["jobs"].items()]
95106

96107
if len(inserts) == 0:
97108
return

mrq/monkey.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,6 @@
1+
from future import standard_library
2+
standard_library.install_aliases()
3+
from past.builtins import basestring
14
from .context import get_current_job, get_current_worker
25
import time
36
import random
@@ -339,7 +342,7 @@ def connect(old_method, self, *args, **kwargs):
339342

340343
return ret
341344

342-
from httplib import HTTPConnection, HTTPSConnection
345+
from http.client import HTTPConnection, HTTPSConnection
343346

344347
patch_method(HTTPConnection, "request", request)
345348
patch_method(HTTPConnection, "connect", connect)
@@ -392,7 +395,7 @@ def _Cursor__send_message(self, *args, **kwargs):
392395
collection = self._Cursor__collection.name # pylint: disable=no-member
393396

394397
if collection == "$cmd":
395-
items = self._Cursor__spec.items() # pylint: disable=no-member
398+
items = list(self._Cursor__spec.items()) # pylint: disable=no-member
396399
if len(items) > 0:
397400
subtype, collection = items[0]
398401

0 commit comments

Comments
 (0)