Wait Click Join in YTsaurus Flow (Java)

The pipeline performs a join of streams of hits and actions within a time window. For each hit, the system waits for related actions (show, click) within a specified time interval, then generates a merged event.

Source code (Java)

Source code (Kotlin)
For a detailed description of the task and join logic, see the C++ version description.

Problem statement

You have a system with three types of events:

  • hit — a user request event. It contains hit_id, hit_time, and hit_payload.
  • action show — an event indicating that the user saw the response. It contains hit_id, hit_time, action_time, and is_click = false.
  • action click — an event indicating that the user clicked the response. It contains hit_id, hit_time, action_time, and is_click = true.

Data arrives in two streams: hit and action.

The output produces a join: for each show, you attach information about whether a click occurred and the hit_payload from the hit. Hits without shows are ignored. The system waits for action events for wait_for_actions seconds after hit_time.

Data flow diagram

JoinProcessFunction

The core pipeline logic is implemented in JoinProcessFunction, which processes two streams — hit and action:

Key patterns

Late data filtering

Messages that arrive with an event time less than the current watermark are considered late and are ignored. For more details about watermarks, see the Watermarks & Timers section.

Timers for closing the window

When you receive a hit, you set a timer for hit_time + wait_for_actions. The timer triggers when the watermark reaches maxTime, meaning the system is sure that all events with event_time < maxTime have been processed.

ExternalStateAccessor with PayloadBuilder

The state is stored in an external dynamic table. PayloadBuilder lets you update individual state fields without full re-serialization:

State cleanup

After the timer fires and the output event is generated, you clear the state using stateAccessor.clear() to avoid accumulating outdated profiles.

Data models

POJO classes with JPA annotations are defined for typed message handling.

Warning

You must ensure that the order of fields in the typed data model matches the order of columns in the stream definition in the static spec. If you break the field order, it can lead to hard-to-diagnose errors.

If the stream is not described in the static spec, the typed model sets the column order. Keep that order compatible with the output queues and tables.

Hit

Action

JoinedAction

Pipeline configuration (Spring)

Source code (Java)

Source code (Kotlin)
The example uses the Spring Boot integration for configuration. The join computation is registered with the @FlowComputation annotation on the JoinProcessFunction class:

Typed streams are declared via ComputationProvider (the getStreams() method):

Key points:

  • The join computation (JoinProcessFunction) is registered with the @FlowComputation(id = "join") annotation.
  • You register three typed streams: hit, action, and joined_action.
  • Typed streams let you use message.getPayload() to retrieve POJO objects.

Entry points

PipelineMain

This single entry point starts the pipeline or serves as its companion, depending on YT_FLOW_MODE.

Static spec

Computation join

"join" = {
    "computation_class_name" = "TJoin";
    "group_by_schema" = [
        {"name" = "hash"; "expression" = "farm_hash(hit_id)"; "type" = "uint64"};
        {"name" = "hit_id"; "type" = "string"};
        {"name" = "hit_time"; "type" = "uint64"};
    ];
    "input_stream_ids" = ["action"; "hit"];
    "output_stream_ids" = ["joined_action"];
    "external_state_managers" = {
        "/join-state" = {
            "external_state_manager_class_name" = "NYT::NFlow::TSimpleExternalStateManager";
            "parameters" = {
                "path" = "//path/to/state";
            };
        };
    };
    "parameters" = {
        "wait_for_actions" = "10s";
    };
    "timers" = {
        "timer" = {};
    };
};
  • group_by_schema is the partitioning key: (hash, hit_id, hit_time).
  • input_stream_ids are two input streams: action and hit.
  • output_stream_ids is one output stream: joined_action.
  • external_state_managers is a top-level section that describes External State (on the same level as parameters). The key ("/join-state") is the state name that starts with /; you pass the same name to the StateDescriptors.external("/join-state") descriptor. external_state_manager_class_name is the registered manager class (NYT::NFlow::TSimpleExternalStateManager for the standard option). parameters/path is the path to the YTsaurus dynamic table; the key column schema of the table must match group_by_schema.
  • parameters.wait_for_actions is the time you wait for action events after a hit. You can access it in Java via RuntimeContext.getComputationParameters().
  • timers.timer declares the timer stream for closing hits.

Testing

Test source code (Java)

Test source code (Kotlin)

You use TestComputationHarness for testing. It’s a test harness that lets you call doProcess without running the full pipeline.

Set up the test environment

Test hit message processing

Test the full join flow

The full test suite (9 scenarios) is in the source code.

Integration testing

You run integration testing for Java pipelines the same way as for C++ pipelines: by launching the full pipeline, including C++ workers, queues, and streams. For details, see Integration tests.

Key ideas

  1. Time window: The system waits for action events for wait_for_actions seconds after hit_time. You discard all events outside this window.

  2. Watermark for late data: Events with event_timestamp < watermark are considered late and are discarded. This inevitably leads to some data loss, but it ensures watermark correctness.

  3. External State for data accumulation: You store intermediate data (hit_payload, show_time, click_time) in External State, which is bound to the key (hit_id, hit_time). In this case, you’re not limited to using only External State; an implementation with Internal State is also possible.

  4. Timer for closing the window: A timer with triggerTimestamp = hit_time + wait_for_actions ensures result generation after the wait window closes.

  5. Idempotency of timers: You set the timer on every message, but with the same triggerTimestamp. Flow deduplicates timers with the same key and triggerTimestamp.

  6. State cleanup: After you generate the result, you clear the state via stateAccessor.clear(), which removes the row from the table.

See also