【问题标题】:Kafka KStream-KStream leftjoin windowed with custom TimestampExtractor cause Skipping record for expired segmentKafka KStream-KStream leftjoin 使用自定义 TimestampExtractor 窗口化导致跳过过期段的记录
【发布时间】:2020-11-02 20:20:45
【问题描述】:

我正在尝试使用自定义 TimestampExtractor 加入 2 KStream (stream1, stream2), 如果我的 2 个事件的时间戳彼此非常接近,我会得到:

[my-app-client-StreamThread-1] WARN org.apache.kafka.streams.state.internals.AbstractRocksDBSegmentedBytesStore - Skipping record for expired segment.

为了测试,我尝试不使用自定义 TimestampExtractor,如果我的生产者足够快地发送事件并遵守我的窗口持续时间配置,它就可以工作。

有什么想法吗?

我查看了文档,在加入 2 个 KStreams 时我没有看到有关自定义 TimestampExtractor 的限制?

以下是有关该问题的更多详细信息:

我的 TimestampExtractor 从事件负载中提取时间戳:

public class EventTimestampExtractor implements TimestampExtractor {
    @Override
    public long extract(final ConsumerRecord<Object, Object> record, final long previousTimestamp) {
        final Event event = (Event) record.value();
        final long timestamp = = event.myTimestamp;
        return timestamp;
    }
}

这是我的申请:

    final Properties props = new Properties();
    props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, EventTimestampExtractor.class);
    
    ...
    
    final StreamsBuilder builder = new StreamsBuilder();


    final KStream<String, Event> stream1 =
            builder.stream("topic-left", Consumed.with(Serdes.String(),
                    EventSerde.serde()));

    final KStream<String, Event>  stream2 =
            builder.stream("topic-right", Consumed.with(Serdes.String(),
                    EventSerde.serde()));


    stream1.leftJoin(stream2,
            (eventLeft, eventRight) -> {

                ... processing ...
                
                Data data = merge(eventLeft.data, eventRight.data);
                
                return data;

            },
            JoinWindows.of(Duration.ofMillis(1000)).grace(Duration.ofMillis(60000)),
            StreamJoined.with(Serdes.String(), EventSerde.serde(), EventSerde.serde())
    )
    .peek((key, data) -> {
        LOG.debug(key + data);
    });

    ...
    
    final KafkaStreams streams = new KafkaStreams(builder.build(), props);
    streams.start();
    

之后,我使用 kafka-console-producer 控制台注入,"topic-left" 中的eventLeft"topic-right" 中的eventRight,其中:

  • eventLeft.myTimestamp = T (in ms)
  • eventRight.myTimestamp = T+200 (in ms)

我的问题是,我没有登录 peek() ,而是得到了:

[my-app-client-StreamThread-1] WARN org.apache.kafka.streams.state.internals.AbstractRocksDBSegmentedBytesStore - Skipping record for expired segment.

当我在EventTimestampExtractor 中显示时间戳值时,一切似乎都正常。

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams


    【解决方案1】:

    问题似乎与此 one 重复

    我没有使用 kafka-console-producer 控制台和硬编码时间戳,而是使用了一个应用程序,它可以在没有硬编码时间戳的情况下生成我的数据并且它可以工作。

    我不明白为什么它不适用于 kafka-console-producer 控制台,在记录内部和提取的时间戳方面肯定存在一些棘手的行为。

    【讨论】:

      猜你喜欢
      • 2018-09-18
      • 1970-01-01
      • 2017-06-02
      • 2017-01-08
      • 1970-01-01
      • 2020-10-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多