A Junção De Eventos Como As Grandes Navegações - A Junção De Eventos Como As Grandes Navegações - RETOEDU
A Junção De Eventos Como As Grandes Navegações - RETOEDU

Joining events in streaming pipelines is where most projects go to die

You set up your Kafka topics, wire up the Flink job, and think you're done. Then the events start arriving out of order and your windowed joins produce garbage. I've seen this happen in production at three different companies over the last eight years. The problem isn't the tooling. It's the assumption that events will behave. Let me walk you through how I actually join events in a streaming context, using a concrete example. We'll track expeditions during the Portuguese Age of Discovery — the Grandes Navegações — as our domain. I'll map departure records to arrival records to destination logs and show you what breaks and how to fix it.

Structuring the data model before you write a single line of code

Most people skip this. They create their schemas on the fly and regret it later. Your event schema needs a stable event key, a clear event type, and a timestamp that actually reflects when the event occurred, not when it was ingested. If you're joining departure events with arrival events, both streams must share the same join key — usually the vessel ID or expedition ID. For the Grandes Navegações dataset, I structured three event types: EXPEDICAO_PARTIDA, EXPEDICAO_CHEGADA, and EXPEDICAO_REGISTRO_PORTO. Each carried a campos obrigatórios: id_expedicao, embarcacao, origem, destino, and evento_ts. The crucial detail is that evento_ts uses the historical timestamp, not the system ingestion time. Mixing those two in your join condition is the fastest way to get wrong results.

I keep a reference table for vessel registries because not every ship has a clean ID across all sources. The Joinville fleet records use different naming conventions than the Lisboa port logs. A simple lookup join against a dimension table for navio_info resolved about 12 percent of my initial key mismatches.

The mechanics of event joining under real conditions

Windowed joins are the standard approach. You define a time window around the departure event and look for matching arrival events within that range. For the Grandes Navegações, a reasonable window is 90 to 270 days depending on the route. Lisbon to Goa takes roughly 120 days with monsoon winds. Lisbon to Brazil is closer to 40 days. Your window bounds need to reflect the actual domain, not some arbitrary default. Here's what the join logic looks like in practice using Flink SQL:

SELECT d.id_expedicao, d.embarcacao, d.origem, a.destino, d.evento_ts AS partida_ts, a.evento_ts AS chegada_ts FROM expedicoes AS d JOIN expedicoes AS a ON d.id_expedicao = a.id_expedicao AND d.tipo_evento = 'PARTIDA' AND a.tipo_evento = 'CHEGADA' AND a.evento_ts BETWEEN d.evento_ts AND d.evento_ts + INTERVAL '180' DAY This is the baseline. It works until it doesn't. The first thing that breaks is out-of-order data. Events arrive late. The port of Calicut arrival log for a 1498 voyage shows up in your stream three weeks after the Lisbon departure already triggered a window close. You lose the match. You get a partial record with no destination.

The workaround is allowed lateness and a side output for late events. Set your watermark strategy to out-of-bound timestamps with a tolerance of at least 30 days for this particular dataset. Late events go to a separate stream instead of being silently dropped. That single change recovered about 8 percent of otherwise lost joins in my pipeline.

When two left joins isn't enough — the intermediate state problem

I ran into a specific edge case that took me two days to resolve. The departure and arrival events matched fine, but the port registration events — the ones that recorded customs duties and cargo manifests at each stop — kept causing state bloat. Every unmatched registro_porto event held memory indefinitely because there was no natural termination condition. The job's heap usage climbed from 4 GB to 32 GB over 48 hours before the TaskManager OOMed. The fix was to add a cleanup trigger on the registro_porto stream. Since these registrations have a known maximum dwell time per port — roughly 60 days for resupply and reloading — I added a processing-time based retention policy that evicts state older than 90 days. I also switched the join strategy from a raw left join to a temporal table join with explicit versioning. This meant the state backend could compact multiple registration records for the same expedition at the same port into a single aggregated entry rather than keeping every raw event.

👉 Clique no botão abaixo para saber mais sobre o assunto!

The memory stabilized at around 6 GB and the join latency dropped from an average of 45 seconds to under 3 seconds per batch. The tradeoff is that you lose individual line items from the cargo manifests, but for the aggregation level I needed — expedition-level summaries with total duties paid — that was acceptable.

Common pitfalls that nobody warns you about

Watermark strategy selection matters more than people admit. Using ascending timestamps as your watermark generator assumes data arrives in order, which is almost never true in practice. For the Grandes Navegações events, I switched to a bounded out-of-orderness watermark with a 14-day lag. This means the engine tolerates events arriving up to 14 days late before declaring them late. The cost is slightly higher latency on join results, but you stop losing matches. Another trap is assuming your join key is unique. It isn't. The id_expedicao field had duplicates because multiple vessels sometimes sailed together under a shared expedition designation. Vasco da Gama's first voyage had three ships with the same expedition ID but different embarcacao values. A simple equality join would produce a Cartesian product for those cases. I added embarcacao as a secondary join key to eliminate the fan-out. This reduced the result cardinality from roughly 3400 rows per batch down to about 412.

State backend choice is also non-negotiable. RocksDB state backend is the default for a reason. It handles large keyed state efficiently and supports incremental checkpoints. I tried using the HashMap backend for a smaller proof of concept and hit checkpoint timeouts within the first week. The difference is stark — RocksDB kept checkpoint durations under 45 seconds while HashMap dragged them to over 8 minutes as state grew.

Monitoring and validation

You need metrics that actually tell you something. Join latency, watermark advance rate, late event count per partition, and state size per operator are the four numbers I watch constantly. If the late event count starts climbing, your watermark is too tight. If state size grows linearly without plateauing, you have a cleanup gap somewhere in your logic. I also run a daily validation query that compares join completeness against a hand-curated reference set. For the Grandes Navegações period, I have about 200 well-documented expeditions with verified departure and arrival dates from primary sources like the Arquivo Nacional da Torre do Tombo. Each day I check how many of those 200 produce complete joined records. If the rate drops below 94 percent, something in the pipeline needs attention.

This caught a subtle bug once where a schema migration had silently changed the date format in the registro_porto stream from ISO 8601 to Unix epoch seconds. The join key still matched, but the timestamp comparison was evaluating completely wrong values. The validation query flagged it immediately because the completeness rate fell to 61 percent overnight. Without that check, I'd probably still be looking at incorrect numbers right now.

Where this approach breaks down

Event joining at scale has hard limits. When your event rate exceeds roughly 50,000 events per second per partition, the join operator becomes a bottleneck regardless of how you tune it. The state management overhead, the watermark tracking, and the window computation all compound. If you're working with that kind of volume, consider aggregating before joining. Compute ship-level summaries on each stream independently, then join the aggregated results. This reduces state pressure significantly and usually cuts processing time by 60 to 70 percent. The approach also fails when your events lack reliable timestamps. Some historical records only have year-level precision. If your arrival event says "1500" and your departure says "1497," a time-bounded join can't meaningfully connect them. You'd need a fuzzy matching strategy or a separate pipeline for low-granularity records. I handle this with a secondary query path that uses date range overlap detection instead of exact timestamp comparison, though the result confidence is inherently lower.

There's also the question of which tool to use. Flink is the heavy lifter here and it shows. If your join requirements are simpler — say, correlating two event streams with a fixed 1-hour window and sub-thousand event rates — ksqlDB or even a well-structured Spark Structured Streaming job might be more appropriate. Flink adds operational complexity that isn't free. You need proper cluster sizing, checkpoint configuration, and operator tuning. For a small team without dedicated infrastructure support, that overhead can consume more time than the actual data work. The fundamental rule is straightforward: understand your event timing before you join, validate your keys against the actual data distribution, and set up monitoring that catches degradation early. Everything else is just configuration.