Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion luigi/contrib/hadoop.py
Original file line number Diff line number Diff line change
Expand Up @@ -957,7 +957,7 @@ def dump(self, directory=""):
if self.__module__ == "__main__":
d = pickle.dumps(self)
module_name = os.path.basename(sys.argv[0]).rsplit(".", 1)[0]
d = d.replace(b"(c__main__", "(c" + module_name)
d = d.replace(b"(c__main__", b"(c" + module_name.encode("utf-8"))
open(file_name, "wb").write(d)

else:
Expand Down
22 changes: 22 additions & 0 deletions test/contrib/hadoop_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import os
import sys
import json
import tempfile
import unittest

import luigi
Expand Down Expand Up @@ -498,3 +499,24 @@ def test_kill_last_application_on_interrupt(self):
]
subprocess = self._run_and_track_with_interrupt(err_lines)
subprocess.call.assert_called_once_with(['yarn', 'application', '-kill', application_id])


class JobTaskDumpTest(unittest.TestCase):
def test_dump_from_main_module(self):
"""`dump` must not raise TypeError for a __main__ job (issue #3284)."""
directory = tempfile.mkdtemp()
job = MyStreamingJob(param='x')
main_module = sys.modules['__main__']
original_module = MyStreamingJob.__module__
original_argv0 = sys.argv[0]
MyStreamingJob.__module__ = '__main__'
setattr(main_module, 'MyStreamingJob', MyStreamingJob)
sys.argv[0] = 'my_job_script.py'
try:
job.dump(directory)
finally:
MyStreamingJob.__module__ = original_module
delattr(main_module, 'MyStreamingJob')
sys.argv[0] = original_argv0

self.assertTrue(os.path.exists(os.path.join(directory, 'job-instance.pickle')))
Loading