Flink Interval Join报Long.MIN_VALUE时间戳错误及业务实现求助
Flink Interval Join 报错:Long.MIN_VALUE timestamp 问题排查与修复
问题背景
使用Flink Interval Join时触发报错:
Long.MIN_VALUE timestamp: Elements used in interval stream joins need to have timestamps meaningful timestamps
业务场景与需求:
- 两个Kafka消费者:A消费
viewTopic,B消费engTopic - 一个Kafka生产者C,输出结果到
delayedTopic - A、B中的事件格式为
view_{id}_{timestamp}和eng_{id}_{timestamp}字符串 - 核心需求:
- 将A中的事件延迟2秒输出
- 基于id匹配,与B中的事件执行Interval Join
- 将关联结果推送到C
- 使用事件时间(Event Time),因Interval Join仅支持事件时间
现有代码实现
主类
public class DataStreamJob { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); Properties properties1 = new Properties(); properties1.setProperty("bootstrap.servers", "localhost:9092"); properties1.setProperty("group.id", "myGroup1"); Properties properties2 = new Properties(); properties2.setProperty("bootstrap.servers", "localhost:9092"); properties2.setProperty("group.id", "myGroup1"); Properties properties3 = new Properties(); properties3.setProperty("bootstrap.servers", "localhost:9092"); // Kafka view事件消费者 FlinkKafkaConsumer<String> viewConsumer = new FlinkKafkaConsumer<>("viewTopic", new SimpleStringSchema(), properties1); DataStream<Event> delayedViewStream = env.addSource(viewConsumer) .map(DataStreamJob::parseEvent) .assignTimestampsAndWatermarks(WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp())) .keyBy(event -> event.getId()) .process(new DelayedEventsProcessFunction()); // Kafka eng事件消费者 FlinkKafkaConsumer<String> engConsumer = new FlinkKafkaConsumer<>("engTopic", new SimpleStringSchema(), properties2); DataStream<Event> directEngStream = env.addSource(engConsumer) .map(DataStreamJob::parseEvent) .assignTimestampsAndWatermarks(WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp())) .keyBy(event -> event.getId()); // Interval Join关联两个流 DataStream<Event> joinedStream = delayedViewStream.keyBy(event -> event.getId()) .intervalJoin(directEngStream.keyBy(event -> event.getId())) .between(Time.seconds(-10), Time.seconds(10)) .process(new EventJoinProcessFunction()); // Kafka结果生产者 KeyedSerializationSchema<String> keyedSerializationSchema = new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema()); FlinkKafkaProducer<String> kafkaProducer = new FlinkKafkaProducer<>( "delayedTopic", keyedSerializationSchema, properties3, FlinkKafkaProducer.Semantic.EXACTLY_ONCE); joinedStream.map(Event::eventToString).addSink(kafkaProducer); env.execute(); }
DelayedEventsProcessFunction(延迟处理逻辑)
public class DelayedEventsProcessFunction extends KeyedProcessFunction<String, Event, Event> { private transient MapState<Long, Event> delayedEventsState; @Override public void open(Configuration parameters) throws Exception { MapStateDescriptor<Long, Event> delayedEventsDescriptor = new MapStateDescriptor<>( "delayedEventsState", Types.LONG, TypeInformation.of(new TypeHint<Event>() {}) ); delayedEventsState = getRuntimeContext().getMapState(delayedEventsDescriptor); } @Override public void processElement(Event event, KeyedProcessFunction<String, Event, Event>.Context ctx, Collector<Event> out) throws Exception { long currentTimestamp = event.getTimestamp(); long delayedTimestamp = currentTimestamp + TimeUnit.SECONDS.toMillis(2); // 注册处理时间定时器 ctx.timerService().registerProcessingTimeTimer(delayedTimestamp); delayedEventsState.put(delayedTimestamp, event); } @Override public void onTimer(long timestamp, OnTimerContext ctx, Collector<Event> out) throws Exception { Event delayedEvent = delayedEventsState.get(timestamp); if (delayedEvent != null) { out.collect(delayedEvent); delayedEventsState.remove(timestamp); } } }
EventJoinProcessFunction(关联逻辑)
public class EventJoinProcessFunction extends ProcessJoinFunction<Event, Event, Event> { @Override public void processElement(Event left, Event right, ProcessJoinFunction<Event, Event, Event>.Context ctx, Collector<Event> out) throws Exception { out.collect(new Event("join:"+left.getType()+right.getType(), left.getId()+right.getId(), left.getTimestamp())); } }
Event类
public class Event { private final String type; private final String id; private long timestamp; public Event(String type, String id, long timestamp) { this.type = type; this.id = id; this.timestamp = timestamp; } public String getType() { return type; } public String getId() { return id; } public long getTimestamp() {return timestamp;} public String eventToString() { return String.format("%s_%s_%s", type, id, timestamp); } }
报错原因分析
报错核心是经过DelayedEventsProcessFunction处理后的流丢失了事件时间和水位线(Watermark)信息:
KeyedProcessFunction不会自动传递上游的时间属性,输出的事件没有有效的事件时间戳- Interval Join依赖事件时间进行关联,无法获取有效时间戳时会触发
Long.MIN_VALUE错误
修复方案
1. 为Event类添加时间戳Setter方法
public class Event { // ... 原有代码 ... public void setTimestamp(long timestamp) { this.timestamp = timestamp; } // ... 原有代码 ... }
2. 修改DelayedEventsProcessFunction,更新事件时间戳
在定时器触发时,将事件的时间戳更新为延迟后的时间:
@Override public void onTimer(long timestamp, OnTimerContext ctx, Collector<Event> out) throws Exception { Event delayedEvent = delayedEventsState.get(timestamp); if (delayedEvent != null) { // 将事件时间戳更新为延迟后的时间 delayedEvent.setTimestamp(timestamp); out.collect(delayedEvent); delayedEventsState.remove(timestamp); } }
3. 为延迟后的流重新分配事件时间与水位线
在process(new DelayedEventsProcessFunction())之后,重新添加assignTimestampsAndWatermarks:
DataStream<Event> delayedViewStream = env.addSource(viewConsumer) .map(DataStreamJob::parseEvent) .assignTimestampsAndWatermarks(WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp())) .keyBy(event -> event.getId()) .process(new DelayedEventsProcessFunction()) // 重新分配事件时间和水位线 .assignTimestampsAndWatermarks(WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp()));
4. 验证Interval Join时间范围
根据延迟逻辑,确认between(Time.seconds(-10), Time.seconds(10))是否符合业务预期,可根据实际需求调整时间区间。
本地环境搭建
- 下载Mac版IntelliJ
- 下载最新版Kafka
- 解压Kafka包:
tar -xzf kafka_2.13-3.4.0.tgz - 进入Kafka目录:
cd kafka_2.13-3.4.0
创建3个Topic
bin/kafka-topics.sh --create --topic viewTopic --bootstrap-server localhost:9092bin/kafka-topics.sh --create --topic engTopic --bootstrap-server localhost:9092bin/kafka-topics.sh --create --topic delayedTopic --bootstrap-server localhost:9092
启动服务与测试
打开5个终端,在Kafka目录下执行:
- 启动ZooKeeper:
bin/zookeeper-server-start.sh config/zookeeper.properties - 启动Kafka Broker:
bin/kafka-server-start.sh config/server.properties - 向viewTopic写入事件:
bin/kafka-console-producer.sh --topic viewTopic --bootstrap-server localhost:9092 < viewEvents.txt - 向engTopic写入事件:
bin/kafka-console-producer.sh --topic engTopic --bootstrap-server localhost:9092 < engEvents.txt - 消费delayedTopic结果:
bin/kafka-console-consumer.sh --topic delayedTopic --bootstrap-server localhost:9092
其中viewEvents.txt和engEvents.txt每行一个事件,格式分别为view_{id}_{timestamp}、eng_{id}_{timestamp},例如view_1_1683787256、eng_1_1683787456
内容的提问来源于stack exchange,提问作者Vivek Chandela
相关产品推荐
相关产品推荐

