The write path: one row per span, one key per trace

Spans are written as they arrive under a random trace id that spreads perfectly across partitions, nothing ever closes a trace, and the tree is rebuilt in the viewer — so a partial trace is the normal steady state, not an error.

Previously

The gateway hands spans to a store one at a time, as they arrive, with no notion of which trace is finished.

Scene 12

One row per span, one key per trace

  1. Watch
  2. Try it
  3. Predict
  4. Capture
a trace is never closed: rows land as they arrivekey = trace id · a trace's rows stay together and the load spreads8 spans / trace6 partitionsretention 14dWRITE TIMEspans arrive out of order, one row at a time0250ms500ms750ms1.0seach row is written on arrival · 8 rows per trace · nothing waits for the trace to endkey = trace id43%p041%p145%p242%p344%p443%p5stored spans (TB) at 1 TB/day — example volume14 / 30retention 14 days · the oldest day expires whole, rows and allREAD TIMEthe tree is rebuilt when somebody asksparent ids are resolved by the reader, not the writernothing is assembled until a read asks for itEach row is written by the process that finished that span. Nothing waits for the others; nothing knows how many others there are.
What to watch for

Where do spans actually go, and who puts a trace back together? Watch two Shopfront requests' spans arrive at the store over 45 seconds — interleaved, out of order, each row written the instant its own process finished that span. Keep your eye on the right-hand half of the canvas while they land.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

TraceStore.write
one row, filed under the key, the moment it arrives
def write(span):
if partition_key == "trace_id":
key = span.trace_id # 128-bit random
else:
key = span.service_name # a few values
partition_of(key).append(row(
trace_id, span_id, parent_span_id,
start_us, duration_us, attributes,
), ttl = retention)
# nothing follows: no assemble(),
# no mark_complete(), no wait_for_siblings()
Viewer.read_trace
group the rows sharing one id, join child to parent
def read_trace(trace_id):
rows = store.get(key = trace_id) # 1 partition
node = {r.span_id: Node(r) for r in rows}
roots = []
for r in rows:
parent = node.get(r.parent_span_id)
if parent: parent.children.append(node[r.span_id])
else: roots.append(node[r.span_id])
return roots
# keyed by service_name, this get() becomes a
# fan-out scan of every partition
TraceStore.expire
the only thing in the system that removes a row
retention = { # all three real
"jaeger on cassandra": 48h, # trace TTL
"tempo blocks": 336h, # block_retention
"aws x-ray": 720h, # fixed, no knob
}
def expire():
for row in partition.scan():
if age(row) > retention:
row.drop() # one row, not one trace
# stored_bytes = write_rate * retention

Where this sits in Build a distributed tracing system (Jaeger / Zipkin style)

Scene 12 of 17, in the Store & read act — Key by trace id, find it, read it, distrust it.. Spans are written as they arrive under a random trace id that spreads perfectly across partitions. Nothing ever closes a trace, so a partial trace is the normal steady state and the tree is rebuilt at read time.

Up next. Reading a trace is one cheap lookup — if you have the id. Nobody has the id. So how do you find the trace that matters?

All 17 scenes in Build a distributed tracing system (Jaeger / Zipkin style) · Every curriculum

Built with Arqly
Every scene in Build a distributed tracing system (Jaeger / Zipkin style) builds on the one before it.All 17 Build a distributed tracing system (Jaeger / Zipkin style) scenes