The problem
Data teams chain batch jobs: fetch raw data, clean it, aggregate it, publish a report. With enough of them, the hard part stops being any single job and becomes the plumbing between them. The Luigi README is blunt about why: you "chain many tasks, automate them, and failures will happen."
What they did
Spotify wrote Luigi, "a Python module that helps you build complex pipelines of batch jobs. It handles dependency resolution, workflow management, visualization etc." In their words: "We use Luigi internally at Spotify to run thousands of tasks every day, organized in complex dependency graphs." Most are Hadoop jobs behind recommendations, toplists, A/B test analysis, external reports and internal dashboards.
They open-sourced it, and the README lists other companies using it, including Foursquare, Stripe, Buffer, SeatGeek and Skyscanner.
The design idea
Two rules make a pipeline safe to re-run after a failure:
- A step whose output already exists is done. Re-running the pipeline skips it instead of redoing hours of work.
- Outputs appear all at once. Luigi makes file-system operations atomic, so, as the README puts it, "your data pipeline will not crash in a state containing partial data."
The part any team can copy
You don't need Luigi for the idea. Here it is in a dozen lines, with a dictionary standing in for files on disk:
store = {} # stands in for output files
runs = [] # which steps actually did work
def step(name, deps, make):
if name in store:
return store[name] # output exists: skip
inputs = [step(*d) for d in deps] # make sure dependencies ran first
runs.append(name)
store[name] = make(*inputs) # publish only a finished result
return store[name]
raw = ("raw", [], lambda: [3, 1, 2])
clean = ("clean", [raw], lambda xs: sorted(xs))
report = ("report", [clean], lambda xs: f"top={xs[-1]}")
print(step(*report), runs)
runs.clear()
print(step(*report), runs) # second run: nothing recomputedtop=3 ['raw', 'clean', 'report'] top=3 []
With real files, "publish only a finished result" means writing to a temporary name and renaming it when complete, so a crash never leaves a half-written file that looks finished.