Computation in YTsaurus Flow (Java)
Note
Use this page to learn the details for working with computations in Java and Kotlin. For general concepts, see the Computation section.
Computation types
In Flow, you have two types of Computation: Swift and Transform. Your choice affects how you ensure exactly-once guarantees and what transformations you can implement.
| Type | Guarantee approach | Use case |
|---|---|---|
Swift |
The transformation code is deterministic and can be rerun if needed | Stateless transformations |
Transform |
The result is always stored in YTsaurus, so no determinism requirements apply | Stateful transformations |
For more on processing guarantees, see the Processing guarantees section.
For Java and Kotlin pipelines, you choose between Swift or Transform by setting computation_class_name in the static spec:
NYT::NFlow::NCompanion::TTransformCompanionComputation— forTransform.NYT::NFlow::NCompanion::TSwiftMapCompanionComputation— forSwift.NYT::NFlow::NCompanion::TTransformOrderedSourceCompanionComputation— forTransformsource.NYT::NFlow::NCompanion::TSwiftOrderedSourceCompanionComputation— forSwiftsource.
Creating a Computation
In Java and Kotlin code, you create a Computation using Computation.Builder and register it in PipelineContext.
var join = Computation.builder()
.setComputationId("join")
.setProcessFunction(new JoinProcessFunction())
.build();
val join = Computation.builder()
.setComputationId("join")
.setProcessFunction(JoinProcessFunction())
.build()
In the static spec, you create a Computation with the same id (in this example, join):
"join" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"group_by_schema" = [
...
];
"input_stream_ids" = [...];
"output_stream_ids" = [...];
"parameters" = {
...
};
"timers" = {};
};
For more on specs, see the Spec, DynamicSpec, and Config section.
Warning
You must provide processFunction (null is not allowed): you don’t register computations without business logic in Java. If you need passthrough, don’t register the computation in Java at all. Instead, specify the C++ passthrough class in computation_class_name in the static spec (see Passthrough Computation).
SourceComputation
SourceComputation is the top node in the pipeline graph that reads data from external sources. On the worker side, it corresponds to TSwiftOrderedSourceComputation or TTransformOrderedSourceComputation.
In Java, SourceComputation extends Computation. Like Computation, it requires the processFunction parameter.
In the static spec, for deterministic processing without custom state, you specify TSwiftOrderedSourceCompanionComputation. If SourceComputation uses internal state or non-deterministic logic, you specify TTransformOrderedSourceCompanionComputation: the worker materializes the output and records it together with the state and the source offset. The internal state key in such a computation is the source partition key.
Parameters
| Parameter | Required | Description |
|---|---|---|
computationId |
Yes | Unique identifier |
processFunction |
Yes | Function for processing messages |
Creating a SourceComputation
var reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(new HitParsingFunction())
.build();
val reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(HitParsingFunction())
.build()
For a passthrough Source, don’t use Java. Specify NYT::NFlow::TSwiftPassthroughOrderedSourceComputation in computation_class_name in the spec and leave the computation unregistered in the Java companion. For more details, see Passthrough Computation.
Interaction with Worker
When you initialize Worker, it requests information about registered Computation and SourceComputation objects from the Java companion. The source computation on the worker side sends input messages to the Java companion, which applies ProcessFunction to them and returns the result.
Process Function
You implement the business logic for data processing with a Process Function. Choose one of these two interfaces: RowFunction or BatchFunction.
Note
Choosing between RowFunction and BatchFunction depends only on your business logic. RowFunction doesn’t add extra processing overhead compared to BatchFunction because Flow internally transfers data in batches.
RowFunction
RowFunction receives messages and timers one at a time. The interface provides two methods:
onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx)— called for each input message.onTimer(Timer timer, OutputCollector output, RuntimeContext ctx)— called when a timer fires.
Example of a stateless function
public class X2Mapper implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
var messageBuilder = ctx.createMessageBuilder("x2_numbers"); //1
Long number = message.get("number", Long.class); //2
messageBuilder.set("number_x2", number * 2); //3
output.addMessage(messageBuilder.finish()); //4
}
}
class X2Mapper : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
val messageBuilder = ctx.createMessageBuilder("x2_numbers") //1
val number: Long? = message.get("number", Long::class.java) //2
messageBuilder.set("number_x2", number!! * 2) //3
output.addMessage(messageBuilder.finish()) //4
}
}
Let’s walk through the code line by line:
ctx.createMessageBuilder("x2_numbers")— you create aMessageBuilderfor the output stream with id =x2_numbers. The stream with this identifier must be present in theoutput_stream_idslist in the computation’s static spec.message.get("number", Long.class)— you get the value of thenumberfield from the incoming message. You must pass the value’s class to theMessage#getmethod to unambiguously convert the serialized form to a Java object.messageBuilder.set("number_x2", number * 2)— you write the value to thenumber_x2field. This field must be present in the schema of thex2_numbersstream in the static spec.output.addMessage(messageBuilder.finish())— thefinishmethod returns the completed message, which you add to theOutputCollector.
BatchFunction
BatchFunction receives the entire list of messages and timers that come from the worker. The interface provides two methods:
onMessages(List<ExtendedMessage> messages, OutputCollector output, RuntimeContext ctx)— called for a batch of messages.onTimers(List<Timer> timers, OutputCollector output, RuntimeContext ctx)— called for a batch of timers.
A batch corresponds to one worker request and may contain messages with different keys; per-key grouping, if needed, is done in user code (see Companion).
Example of a batch function
public class X2BatchMapper implements BatchFunction {
@Override
public void onMessages(List<ExtendedMessage> messages, OutputCollector output, RuntimeContext ctx) {
var messageBuilder = ctx.createMessageBuilder("x2_numbers"); //1
for (var message : messages) { //2
Long number = message.get("number", Long.class); //3
messageBuilder.set("number_x2", number * 2); //4
output.addMessage(messageBuilder.finish()); //5
}
}
}
class X2BatchMapper : BatchFunction {
override fun onMessages(messages: List<ExtendedMessage>, output: OutputCollector, ctx: RuntimeContext) {
val messageBuilder = ctx.createMessageBuilder("x2_numbers") //1
for (message in messages) { //2
val number: Long? = message.get("number", Long::class.java) //3
messageBuilder.set("number_x2", number!! * 2) //4
output.addMessage(messageBuilder.finish()) //5
}
}
}
The key differences from RowFunction are:
- You create
MessageBuilderonce for the entire batch (line 1). - The
finish()method returns the completed message and resetsMessageBuilderto its initial state, so you can reuse it for the next message (line 5).
Registering in PipelineContext
You must register all Computation objects and typed streams (created via FlowStreams.typed) in PipelineContext before you run GrpcServerExecution.
You don’t need to register untyped streams (created via FlowStreams.raw). Flow creates them automatically based on the streams block in the static spec.
Learn more about Typed Streams.
var context = new PipelineContext();
// Register Computation objects.
Computation join = Computation.builder()
.setComputationId("join")
.setProcessFunction(new JoinProcessFunction())
.build();
context.registerComputation(join);
SourceComputation reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(new HitParsingFunction())
.build();
context.registerComputation(reader);
// Register typed streams.
context.registerStream(FlowStreams.typed("hit", Hit.class));
context.registerStream(FlowStreams.typed("action", Action.class));
context.registerStream(FlowStreams.typed("joined_action", JoinedAction.class));
val context = PipelineContext()
// Register Computation objects.
val join: Computation = Computation.builder()
.setComputationId("join")
.setProcessFunction(JoinProcessFunction())
.build()
context.registerComputation(join)
val reader: SourceComputation = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(HitParsingFunction())
.build()
context.registerComputation(reader)
// Register typed streams.
context.registerStream(FlowStreams.typed("hit", Hit::class.java))
context.registerStream(FlowStreams.typed("action", Action::class.java))
context.registerStream(FlowStreams.typed("joined_action", JoinedAction::class.java))
Warning
Each Computation and stream must have a unique ID that matches the IDs in the static spec. If you try to register a Computation or stream with an ID that already exists, you’ll get an error and won’t be able to start the companion.
RuntimeContext
RuntimeContext gives you access to the computation’s execution context. Key methods:
| Method | Description |
|---|---|
ctx.createMessageBuilder(streamId) |
Create a MessageBuilder for the specified output stream |
ctx.getComputationParameters() |
Get the computation’s parameters from the spec |
ctx.getEpochInputEventWatermark() |
Get the current watermark for the epoch |
ctx.getProtoStateAccessor(name, message, Class) |
Get the state as a protobuf object linked to the message’s key |
ctx.getYsonStateAccessor(name, message, Class) |
Get the YSON state linked to the message’s key |
ctx.getStateAccessor(name, message, Class, ser, deser) |
Get the state with custom serialization/deserialization |
ctx.getRawStateAccessor(name, message) |
Get the state as a byte array without interpretation |
ctx.getNoOpStateAccessor(name, message) |
Get the state that only stores the presence fact (no value) |
ctx.getExternalStateAccessor(name, message) |
Get the external state linked to the message’s key |
Learn more about working with states in the Working with States (Java) section.
OutputCollector
Use OutputCollector to send processing results:
| Method | Description |
|---|---|
output.addMessage(message) |
Add an output message |
output.addMessage(message, options) |
Add a message with AddMessageOptions controlling distribution and the Swift message ID suffix |
output.addTimer(triggerTimestamp) |
Add a timer with the specified trigger time (eventTimestamp = 0) |
output.addTimer(triggerTimestamp, eventTimestamp) |
Add a timer with the specified trigger time and event time |
output.addTimer(timerStreamId, triggerTimestamp, eventTimestamp) |
Add a timer for a specific timer stream |
output.setParentIds(parentIds) |
Set the parent ID to track the lineage of messages. Returns a new OutputCollector |
For a Swift computation, MessageIdSuffix.payloadHash() or MessageIdSuffix.userDefined(value) makes message identity independent of emission order:
output.addMessage(
message,
AddMessageOptions.builder()
.setMessageIdSuffix(MessageIdSuffix.payloadHash())
.build());
Default options use the message sequence number. Payload-hash and user-defined suffixes are supported only by Swift computations.
Spring Boot
When you use Spring Boot, you register the computation with the @FlowComputation annotation (or @FlowSourceComputation for a source) directly on the ProcessFunction class. The annotation is meta-annotated with @Component, so the class automatically becomes a Spring bean:
@FlowComputation(id = "mapper")
public class WordCountMapper implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
// process the message
}
}
@FlowComputation(id = "mapper")
class WordCountMapper : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
// process the message
}
}
Declare streams as Spring beans of type FlowStream<?> (or via ComputationProvider.getStreams()):
@Configuration
public class StreamConfiguration {
@Bean
public FlowStream<Word> wordsStream() {
return FlowStreams.typed("words", Word.class);
}
}
@Configuration
class StreamConfiguration {
@Bean
fun wordsStream(): FlowStream<Word> = FlowStreams.typed("words", Word::class.java)
}
FlowStreams.typed(...) creates a typed stream that automatically serializes and deserializes messages into Java objects of the specified type. Learn more in the Typed Streams section.
Learn more about registration via annotations and ComputationProvider in the Spring Boot Integration section.
CompanionManager Resource Configuration
To run a companion in Java or Kotlin, you must declare the CompanionManager resource in the static spec:
"CompanionManager" = {
"resource_class_name" = "NYT::NFlow::NCompanion::TJavaCompanionManager";
"parameters" = {
"timeout" = "10s";
"jdk_bin_path" = "/app/ytflow/jdk/bin/java";
"main_class" = "tech.ytsaurus.flow.examples.waitclickjoin.PipelineMain";
"classpath" = "/app/ytflow/lib/*";
};
"dependencies" = {};
};
The resource_class_name parameter specifies the resource class that will run the companion.
For a Java or Kotlin companion, resource_class_name must always be NYT::NFlow::NCompanion::TJavaCompanionManager (it supports both languages via the JVM).
Learn more about the spec in the Spec, DynamicSpec and Config section.