I am currently trying to write some tests for our flink stream processing, and have something such as
env.fromCollection(events)
.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<>(Time.milliseconds(1000)) {
@Override
public long extractTimestamp(CallEvent element) {
return element.getEventTimeStamp();
}
})
.keyBy(CALL_EV_KEY_SELECTOR)
.process(new CallEventKeyedProcessFunction())
.map(new CallTimedEntityMapper())
.addSink(sink);
env.execute();
in CallEventKeyedProcessFunction, when attempting to call ctx.timerService().currentWatermark();, this will always return Long.MIN. Is there anything obvious I might be missing?