TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Version the Change and Keep Both Sides Alive

Continue in TT Lab

Goal

You build a contract tool contract.py that handles JSON Lines with versioned schemas. You classify the changes between versions by kind, judge backward compatibility and forward compatibility according to Avro's schema resolution rules, and get through the transition period in which several versions flow mixed together with a single read schema.

Why it matters

The writer and the reader do not change at the same moment. Between them there is always a period when only one side has the new code, and data keeps flowing during that period. So what you should ask when designing a schema change is not "is this change correct" but "does it survive being deployed in either order." The answer splits into three things. Adding a field with a default is safe in both directions, adding a required field without a default makes new reading code unable to read old data, and in type changes, the widening direction and the narrowing direction give exactly opposite results. A rename is, looking only at the schema, one addition and one deletion, and it comes back to a single event only through an alias. But because only the reader schema's aliases are used, a rename survives in only one direction. There is only one way to get through the transition period. You attach generous defaults to the read schema so that even the old versions can be read. This lab shows in numbers how many rows that one default saves. The grader does not trust your text. It sets up the schemas and lines that the grader created in a temporary directory, actually runs your tool, and compares the classification and judgments with values that the grader implemented separately. The field names, types, and amounts change on every run.

Steps

  1. Create and run /root/evolve/gen_stream.py to create v1.json through v5.json in /root/evolve/schemas, and in /root/evolve/stream v1.jsonl through v5.jsonl with 12 lines per version, and mixed.jsonl with the five versions mixed.
  2. In /root/evolve/contract.py, create fields <스키마> so that it produces the version number, the field names, required or optional, and the defaults.
  3. Add diff <옛 스키마> <새 스키마> to classify additions and deletions. Additions are split into those with a default and those without.
  4. Make diff look at the new version's aliases and group a rename into a single event. The grouped names drop out of the addition and deletion lists.
  5. Make diff split type changes into widening and narrowing. If it is in the promotion table, it is widening, and otherwise narrowing.
  6. Add compat <옛 스키마> <새 스키마> to judge backward and forward compatibility and produce the reasons as fixed codes.
  7. Add read <읽기 스키마> <스키마폴더> <파일> to read lines from mixed versions as a single version, create a tolerant read schema /root/evolve/reader.json, and write the difference between the two reads in /root/evolve/window.json.
  8. Leave the history of the five versions in /root/evolve/evolve_report.json and /root/evolve/evolve_report.md.

Reference

Create five versions and emit them

Create and run /root/evolve/gen_stream.py to create v1.json through v5.json in /root/evolve/schemas, and in /root/evolve/stream v1.jsonl through v5.jsonl, and mixed.jsonl with the five versions mixed.

If you make only one thing differ per version, you can later see what the judgment catches. v2 adds a field with a default, v3 adds a field without a default, v4 renames (attach an alias) while widening a type, and v5 narrows that type again. Write which version each line is in _v.

Separate required from optional by default

In /root/evolve/contract.py, create fields <스키마> (the placeholder stands for the schema) so that it produces version, names, required, optional, and defaults as JSON. A field with a default is optional and one without is required.

Produce names in the declaration order as is, and required and optional sorted. You only need to look at whether the default key exists, and even if the default is an empty string or 0, it is optional — you decide by the existence of the key, not the value.

Separate what was added from what was removed

Add diff <옛 스키마> <새 스키마> (the placeholders stand for the old and new schemas) to produce the three lists added_with_default, added_required, and removed. Produce all three lists sorted.

You only need the difference of the name sets. Split the added fields into two branches by whether the default key exists in the new version's declaration. In this step you do not yet need to look at aliases and type changes.

Group a rename into a single event

Add renamed to the diff response. If the aliases of a field that exists only in the new version contains a name that exists only in the old version, the two are the same field. The grouped names drop out of added_* and removed.

A rename is, looking only at the schema, one addition and one deletion. The alias is the only device that turns the two back into a single event. Do not guess from value samples — here we are the side that writes the aliases, so there is no reason to guess.

Separate widening from narrowing

Add widened and narrowed to the diff response. An entry is [새이름, 옛타입, 새타입] (new name, old type, new type), and if the old type is promoted to the new type, it is widening, and otherwise narrowing. It also looks at fields whose type changed at the same time as the name changed.

If you write the promotion table as a single dictionary, the judgment becomes one line. int goes to long, float, and double, long goes to float and double, and float goes to double. The same type must also count as a promotion so that the compatibility judgment later is simple.

Judge backward and forward separately

Add compat <옛 스키마> <새 스키마> (the placeholders stand for the old and new schemas) to produce backward, forward, and reasons. Use only the four codes in the reference section for the reasons and produce them sorted.

Build one function that checks whether the reader can read the writer's data, and swap its arguments to make the two directions. For backward compatibility the reader is the new version, and for forward compatibility the reader is the old version. The fact that aliases are used only from the reader's schema is what decides the result here.

Get through the transition period with one read schema

Add read <읽기 스키마> <스키마폴더> <파일> (the placeholders stand for the read schema, the schema folder, and the file), and create a tolerant read schema /root/evolve/reader.json that reads all five versions. Write the result of reading with schemas/v5.json and the result of reading with reader.json in /root/evolve/window.json as strict, tolerant, and amount_total.

A strict read schema throws away the old versions wholesale because of a field without a default. How many rows come back to life when you attach a default to that field is the answer of this step. The alias must also be on the reader side for it to follow the old name. Write the tolerant side's value for amount_total.

Leave the version history on one page

Run compat on each adjacent pair of versions and write versions, steps, full, and broken in /root/evolve/evolve_report.json, and in /root/evolve/evolve_report.md, write four sections: ## 어떤 판이 있나 ## 어느 방향이 깨지나 ## 전환 기간을 어떻게 넘기나 ## 다음 판에 지킬 것 (in order: which versions exist, which direction breaks, how to get through the transition period, and what to keep in the next version).

steps is a list of {"from": 정수, "to": 정수, "backward": 참거짓, "forward": 참거짓, "reasons": [...]} (integers and booleans). full holds the pairs where both directions work, and broken holds the pairs where even one direction breaks, as [옛판, 새판] (old version, new version). In the report, write the number of rows that came back to life during the transition period as a number.