【问题标题】:Maintain separate KTable维护单独的 KTable
【发布时间】:2020-10-07 05:28:28
【问题描述】:
我有一个主题,其中包含每个会话的用户连接和断开连接事件。我想使用 Kafka 流来处理这个主题并根据某些条件更新 KTable。每条记录都不能更新 KTable。所以我需要处理多条记录才能知道是否需要更新 KTable。
例如,按用户然后按 sessionid 处理流和聚合。如果该用户的至少一个 sessionid 仅有 Connected 事件,则 KTable 必须以在线用户身份更新(如果尚未更新)。
如果用户的所有 sessionId 都有 Disconnected 事件,KTable 必须更新为用户离线,如果还没有。
如何实现这样的逻辑?
我们能否在所有应用实例中实现这个 KTable,以便每个实例在本地都有这些数据?
【问题讨论】:
标签:
apache-kafka
apache-kafka-streams
ktable
【解决方案1】:
听起来是一个相当复杂的场景。
也许,在这种情况下最好使用处理器 API? KTable 基本上只是一个 KV 存储,使用处理器 API,您可以应用复杂的处理来决定是否要更新状态存储。 KTable 本身不允许你应用复杂的逻辑,但它会应用它收到的每个更新。
因此,使用 DSL,您需要进行一些处理,如果您想更新 KTable,请仅针对这种情况发送更新记录。像这样的:
KStream stream = builder.stream("input-topic");
// apply your processing and write an update record into `updates` when necessary
KStream updates = stream...
KTable table = updates.toTable();