Flink State示例
DataStreamSource<Order> sourceStream1 = env.addSource(consumer);
KeyedStream<Order, String> stream1 = sourceStream1.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<Order>(Time.seconds(5)) {
@Override
public long extractTimestamp(Order element) {
return Order.getTime;
}
}).keyBy(Order::getOrderId);
DataStreamSource<Order> sourceStream2 = env.addSource(consumer);
KeyedStream<Order, String> stream2 = sourceStream1.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<Order>(Time.seconds(5)) {
@Override
public long extractTimestamp(Order element) {
return Order.getTime;
}
}).keyBy(Order::getOrderId);
OutputTag<Order> outputTag1 = new OutputTag<>("stream1");
OutputTag<Order> outputTag2 = new OutputTag<>("stream2");
做双流connect
stream1.connect(stream2).process(new CoProcessFunction<Order, Order, Tuple2<Order, Order>>() {
ValueState<Order> state1;
ValueState<Order> state2;
ValueState<Long> timeState;
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
state1 = getRuntimeContext().getState(new ValueStateDescriptor<>("state1", Order.class));
state2 = getRuntimeContext().getState(new ValueStateDescriptor<>("state2", Order.class));
timeState = getRuntimeContext().getState(new ValueStateDescriptor<>("timeState", Long.class));
}
@Override
public void processElement1(Order value, Context ctx, Collector<Tuple2<Order, Order>> out) throws Exception {
Order value2 = state2.value();
if (value2 != null) {
out.collect(Tuple2.of(value, value2));
state2.clear();
ctx.timerService().deleteEventTimeTimer(timeState.value());
timeState.clear();
} else {
state1.update(value);
long time = value.getTime() + 60000;
timeState.update(time);
ctx.timerService().registerEventTimeTimer(time);
}
}
@Override
public void processElement2(Order value, Context ctx, Collector<Tuple2<Order, Order>> out) throws Exception {
Order value1 = state1.value();
if (value1 != null) {
out.collect(Tuple2.of(value1, value));
state1.clear();
ctx.timerService().deleteEventTimeTimer(timeState.value());
timeState.clear();
} else {
state2.update(value);
long time = value.getTime()+ 60000;
timeState.update(time);
ctx.timerService().registerEventTimeTimer(time);
}
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Tuple2<Order, Order>> out) throws Exception {
super.onTimer(timestamp, ctx, out);
if (state1.value() != null) {
ctx.output(outputTag1, state1.value());
}
if (state2.value() != null) {
ctx.output(outputTag2, state2.value());
}
state1.clear();
state2.clear();
}
});
|