Millet Porridge

English version of https://corvo.myseu.cn

0%

Pitfalls of rq in the Python2-to-3 Migration

Preface

Recently I’ve been converting a project from Python2 to Python3. Most of the project’s code was already made as compatible as possible, but for two days I still suffered over upgrading the rq version. I’ll record the problems met during the upgrade; if any friends need it, feel free to reuse this code.

Python2 is expected to stop being supported in 2020 this year, yet a project I maintain is still Python2 code. By the book, upgrading to Python3 basically has no KPI and is a thankless task. But honestly, I’m already quite zen at the company and basically don’t care about KPIs. After I brought it up in a meeting, my boss didn’t object to the upgrade, so I’m tidying up this project — after all, I’ll be maintaining it for a long time. If more problems appear later I’ll record them in the blog too. Maybe the rq problem is a beginning, or maybe an ending.

rq

rq is a Python asynchronous task library that dispatches tasks through redis queues and hands them to workers. Our project depends heavily on rq — switching to another task queue is impossible; the only option is upgrading after writing good unit tests. The upgrade covers two aspects:

  1. upgrading the rq version
  2. upgrading Python

Below I introduce the problems and compatibility solutions for each aspect. The original-to-target versions are:

Item Original version Target version
Python 2.7 3.7
redis 2.10.5 3.3.8
rq 0.6.0 1.11.0

Problems Encountered Upgrading the rq Version

unicode

In Python2 only strings explicitly marked by the user (u'') are unicode strings. In 0.6.0, rq hadn’t even added handling of Chinese in exceptions; since 0.8.0 it added _get_safe_exception_string. Python3 is friendlier to unicode, so this upgrade was relatively smooth.

1
2
3
4
5
6
7
8
9
10
11
12
# This is a modification re-implementing _get_safe_exception_string in user code; if your version is newer than 0.8.0,
# there's no need to add this part.
def _get_safe_exception_string(self, exc_strings):
"""Ensure list of exception strings is decoded and
joined as one string safely.
"""
if rq.__version__ < '1.0':
exc_strings = map(
lambda exc: exc if PY3 else exc.decode("utf-8"), exc_strings)
return ''.join(exc_strings)
else:
return Worker._get_safe_exception_string(exc_strings)

move_to_failed_queue

After rq 1.0, move_to_failed_queue was removed; failed jobs are added directly, no longer requiring manual specification by the user.

1
2
3
4
5
6
7
8
9
10
11
12
13
if rq.__version__ < '1.0':
self.ori_move_to_failed_queue = self.move_to_failed_queue
self.ori_handle_exception = self.handle_exception

kwargs['exception_handlers'] = [
self._handle_job_exception,
self.move_to_failed_queue,
]
else:
# In versions after 1.0, only _handle_job_exception is kept, for printing errors.
kwargs['exception_handlers'] = [
self._handle_job_exception,
]

Job Data Storage and Loading

After rq > 0.13, the exc_info and data in Job are compressed with zlib — that is to say, new-version data cannot be read by old-version rq. A compatibility layer is needed in the custom fetch function:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
try:
raw_data = obj['data']
except KeyError:
raise NoSuchJobError('Unexpected job format: {0}'.format(obj))
try:
self.data = zlib.decompress(raw_data)
except zlib.error:
# Fallback to uncompressed string
self.data = raw_data

raw_exc_info = obj.get('exc_info')
if raw_exc_info:
try:
self.exc_info = as_text(zlib.decompress(raw_exc_info))
except zlib.error:
# Fallback to uncompressed string
self.exc_info = as_text(raw_exc_info)

The complete fetch function can be seen in the appendix.

Time Formats

rq’s time format changed several times; upgrading directly will definitely cause problems. Our custom Job.fetch needs a compatibility layer:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
def utcparse(string):
_d = None
exists_time_format = [
'%Y-%m-%dT%H:%M:%S.%fZ',
'%y-%m-%dt%h:%m:%sz', # rq < 0.9 date format
'%y-%m-%dt%h:%m:%s.%f+00:00', # rq < 0.4 datetime format
]
for t_format in exists_time_format:
try:
_d = datetime.datetime.strptime(string, t_format)
except ValueError:
_d = None
else:
break

if not _d:
log.warning("The %s can't be parsed", string)

return _d

def to_date(date_str):
if date_str is None:
return
else:
return utcparse(as_text(date_str))

The complete fetch function can be seen in the appendix.

Problems the Python Upgrade Caused for rq

The pickle Protocol

rq uses pickle in the following way. The usage itself is fine, but the variable pickle.HIGHEST_PROTOCOL differs between Python2 and Python3. Data dumped in Python3 cannot be read in Python2 because its protocol is older — that is to say, tasks dispatched by rq cannot be read.

1
2
# see rq/job.py
dumps = partial(pickle.dumps, protocol=pickle.HIGHEST_PROTOCOL)

Therefore, in Python3 we need to convert for this situation. The direct method is replacing this dumps function. The code below must be placed at the beginning of the program, and can only be deleted after all our workers are updated to Python3.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
################################################
# This config is for rq. cause in rq/job.py always use the higest protocal
# dumps = partial(pickle.dumps, protocol=pickle.HIGHEST_PROTOCOL)
#
# After refer to https://github.com/rq/rq/issues/598, we limit the
# rq pickle protocol version.
# TODO: Remove this after upgrade python and rq all
if PY3:
try:
import cPickle as pickle
except ImportError: # noqa # pragma: no cover
import pickle
# pickle.HIGHEST_PROTOCOL=2

from unittest import mock
from functools import partial
mydumps = partial(pickle.dumps, protocol=2)
mock.patch('rq.job.dumps', side_effect=mydumps).start()
################################################

Summary

With these changes, many places in the base libraries now carry checks for the Python version or rq version. Problems caused by differing versions can only be solved by upgrading to a unified version. I hope I can fully upgrade our application soon and never be troubled by Python2 again.

Appendix

The mock_fetch function

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
# Manually replace the Job.fetch function in code; put this at the top of the code
if rq.__version__ == '0.6.0':
from ...you_mock import manual_fetch
old_fetch = Job.fetch
@classmethod
def my_fetch(cls, _id, connection=None):
try:
_j = old_fetch(_id, connection)
return _j
except ValueError as e:
log.warning("Please update the rq version {}".format(e))
return manual_fetch(_id, connection=connection)
except Exception as e:
log.error(u"Get job error with: {} {}".format(
e, traceback.format_exc()))
Job.fetch = my_fetc

def manual_fetch(job_id, connection=None):
# type: (int) -> Job
"""manual_fetch
TODO: Remove when rq upgrade, this is only needed in rq==0.6

Give another try:
In differenty versions, rq change time format many times.
like: https://github.com/rq/rq/issues/721
if the fetch not worke, we give a try to parse
the job, hope it can works.


:param job_id: [description]
:type job_id: [type]
"""
_job = Job(job_id, connection=connection)

# patch refresh function
def new_refresh(self): # noqa
"""Overwrite the current instance's properties with the values in the
corresponding Redis key.

Will raise a NoSuchJobError if no corresponding Redis key exists.
"""
from rq.job import decode_redis_hash, unpickle
from rq.utils import as_text
key = self.key
obj = decode_redis_hash(self.connection.hgetall(key))
if not obj:
raise NoSuchJobError('No such job: {0}'.format(key))

def utcparse(string):
_d = None
exists_time_format = [
'%Y-%m-%dT%H:%M:%S.%fZ',
'%y-%m-%dt%h:%m:%sz', # rq < 0.9 date format
'%y-%m-%dt%h:%m:%s.%f+00:00', # rq < 0.4 datetime format
]
for t_format in exists_time_format:
try:
_d = datetime.datetime.strptime(string, t_format)
except ValueError:
_d = None
else:
break

if not _d:
log.warning("The %s can't be parsed", string)

return _d

def to_date(date_str):
if date_str is None:
return
else:
return utcparse(as_text(date_str))

# rq > v0.13 the exc_info and data has been compressed by zlib
# please refer to:
# https://github.com/rq/rq/commit/f500186f3dc652278786a6223e5f24303ad81336
try:
raw_data = obj['data']
except KeyError:
raise NoSuchJobError('Unexpected job format: {0}'.format(obj))

try:
self.data = zlib.decompress(raw_data)
except zlib.error:
# Fallback to uncompressed string
self.data = raw_data

raw_exc_info = obj.get('exc_info')
if raw_exc_info:
try:
self.exc_info = as_text(zlib.decompress(raw_exc_info))
except zlib.error:
# Fallback to uncompressed string
self.exc_info = as_text(raw_exc_info)

self.created_at = to_date(as_text(obj.get('created_at')))
self.origin = as_text(obj.get('origin'))
self.description = as_text(obj.get('description'))
self.enqueued_at = to_date(as_text(obj.get('enqueued_at')))
self.started_at = to_date(as_text(obj.get('started_at')))
self.ended_at = to_date(as_text(obj.get('ended_at')))
self._result = unpickle(obj.get('result')) if obj.get('result') else None # noqa
# self.exc_info = as_text(obj.get('exc_info'))
self.timeout = int(obj.get('timeout')) if obj.get('timeout') else None
self.result_ttl = int(obj.get('result_ttl')) if obj.get('result_ttl') else None # noqa
self._status = as_text(obj.get('status') if obj.get('status') else None)
self._dependency_id = as_text(obj.get('dependency_id', None))
self.ttl = int(obj.get('ttl')) if obj.get('ttl') else None
self.meta = unpickle(obj.get('meta')) if obj.get('meta') else {}

new_refresh(_job)

return _job