【发布时间】: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