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 (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, andhit_payload. - action show — an event indicating that the user saw the response. It contains
hit_id,hit_time,action_time, andis_click = false. - action click — an event indicating that the user clicked the response. It contains
hit_id,hit_time,action_time, andis_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 (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
joincomputation (JoinProcessFunction) is registered with the@FlowComputation(id = "join")annotation. - You register three typed streams:
hit,action, andjoined_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_schemais the partitioning key:(hash, hit_id, hit_time).input_stream_idsare two input streams:actionandhit.output_stream_idsis one output stream:joined_action.external_state_managersis a top-level section that describes External State (on the same level asparameters). The key ("/join-state") is the state name that starts with/; you pass the same name to theStateDescriptors.external("/join-state")descriptor.external_state_manager_class_nameis the registered manager class (NYT::NFlow::TSimpleExternalStateManagerfor the standard option).parameters/pathis the path to the YTsaurus dynamic table; the key column schema of the table must matchgroup_by_schema.parameters.wait_for_actionsis the time you wait for action events after a hit. You can access it in Java viaRuntimeContext.getComputationParameters().timers.timerdeclares the timer stream for closing hits.
Testing
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
-
Time window: The system waits for action events for
wait_for_actionsseconds afterhit_time. You discard all events outside this window. -
Watermark for late data: Events with
event_timestamp < watermarkare considered late and are discarded. This inevitably leads to some data loss, but it ensures watermark correctness. -
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. -
Timer for closing the window: A timer with
triggerTimestamp = hit_time + wait_for_actionsensures result generation after the wait window closes. -
Idempotency of timers: You set the timer on every message, but with the same
triggerTimestamp. Flow deduplicates timers with the same key andtriggerTimestamp. -
State cleanup: After you generate the result, you clear the state via
stateAccessor.clear(), which removes the row from the table.