Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Stateful stream joins: bound waiting time and unmatched records

Last updated: 5 Oct 20265 min read
tutorial
IntermediateBy AITrove Editorial

A stream join keeps state while waiting for related events, so its key, time bounds and expiry behavior are part of correctness.

Choose the join identity

Receipt-submitted and receipt-reviewed events should join by receipt ID. If two review events arrive, define whether the latest version supersedes the first or both represent separate actions. A consumer that joins on customer ID can accidentally match a review to a different receipt. Join cardinality must be stated before streaming state is built.

Bound the waiting period

A review can occur hours after submission. Keep the submission state long enough for normal cases, then route unmatched reviews and expired submissions to separate outputs. Never hold every unmatched key forever; memory grows with unclosed cases. A tight timeout saves memory but can drop valid slow reviews, so measure the delay distribution and document late corrections.

Mind event ordering

The review event may arrive before the submission event even if it occurred later. A processing-order join that discards unmatched reviews will lose it. Event-time ordering and a bounded buffer can help, but the watermark and grace policy still govern when the pair is final. Late routing must include both sides of the join.

Exercise state transitions

Send review before submission for receipt R47, then complete the pair. Send a submission that never gets a review and confirm it expires into a pending output. Send duplicate review IDs and verify no extra pair appears. A backfill replay should produce the same matched, pending and correction totals.

Implementation

python
def pair_receipt_events(submissions, reviews):
    matches = []
    for receipt_id, submission in submissions.items():
        review = reviews.get(receipt_id)
        if review is not None:
            matches.append((receipt_id, submission, review))
    return matches

Performance and operating cost

A batch-style fixture scan is O(S) expected time and O(M) output space for S submissions and M matches. A live join also keeps O(U) state for U unmatched IDs and must expire it under a measured lateness policy.

Common Mistakes

  • Do not join on a broader key than the business unit.
  • Do not hold unmatched state indefinitely.
  • Do not assume arrival order matches occurrence order.

Read next

ai-data
streaming-analytics
Storage details