【问题标题】:Kafka Streams Windowing with Custom TimestampExtractorKafka Streams Windowing with Custom TimestampExtractor
【发布时间】:2018-10-11 11:30:28
【问题描述】:

我正在尝试创建一个 Kafka Streams 应用程序,并尝试在一个时间窗口内计算每个平台的唯一设备。

事件类

public class Event {
    private String eventId;
    private String deviceId;
    private String platform;
    private ZonedDateTime createdAt;
}

我需要时间窗口尊重事件的 createdAt 所以我写了一个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 ZonedDateTime eventCreationTime = event.getCreatedAt();
        final long timestamp = eventCreationTime.toEpochSecond();

        log.trace("Event ({}) yielded timestamp: {}", event.getEventId(), timestamp);

        return timestamp;
    }
}

最后,这是我的流媒体应用代码:

final KStream<String, Event> eventStream = builder.stream("events_ingestion");

eventStream
    .selectKey((key, event) -> {
        final String platform = event.getPlatform();
        final String deviceId = event.getDeviceId());

        return String.join("::", platform, deviceId);
    })
    .groupByKey()
    .windowedBy(TimeWindows.of(TimeUnit.MINUTES.toMillis(15)))
    .count(Materialized.as(COUNT_STORE));

当我将事件推送到event_ingestion 主题时,我可以看到时间戳已记录到应用程序日志中,并且数据正在写入计数存储中。

当我遍历计数存储时,我看到以下内容:

Key: [ANDROID::1@1539000000/1539900000], Value: 2

虽然我的时间窗口是 15 分钟,但关键跨度为 10 天。如果我从流配置中删除 TimestampExtractor 实现(因此返回处理时间),则密钥按预期跨越 15 分钟:

Key: [ANDROID::1@1539256500000/1539257400000], Value: 1

我在这里做错了什么?有什么想法吗?

【问题讨论】:

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


    【解决方案1】:

    TimestampExtractor 使用纪元毫秒值进行窗口化。您正在计算“秒”,这会将消息放入错误的时间窗口。

    【讨论】:

      猜你喜欢
      • 2018-07-03
      • 2017-01-24
      • 1970-01-01
      • 1970-01-01
      • 2019-03-30
      • 2020-04-25
      • 2018-01-19
      • 2019-01-31
      • 1970-01-01
      相关资源
      最近更新 更多