You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink Interval Join报Long.MIN_VALUE时间戳错误及业务实现求助

问题背景

使用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}字符串
  • 核心需求:
    1. 将A中的事件延迟2秒输出
    2. 基于id匹配,与B中的事件执行Interval Join
    3. 将关联结果推送到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:9092
  • bin/kafka-topics.sh --create --topic engTopic --bootstrap-server localhost:9092
  • bin/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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 15:32:31