【发布时间】:2020-05-28 13:04:21
【问题描述】:
我想让两个 faust 代理监听同一个 kafka 主题,但是每个代理在处理事件之前都使用自己的过滤器,并且它们的事件集不会相交。
在文档中我们有一个例子: https://faust.readthedocs.io/en/latest/userguide/streams.html#id4
如果两个代理使用订阅同一主题的流:
topic = app.topic('orders') @app.agent(topic) async def processA(stream): async for value in stream: print(f'A: {value}') @app.agent(topic) async def processB(stream): async for value in stream: print(f'B: {value}')售票员将转发收到的“订单”上的每条消息 两个代理的主题,每当增加引用计数 它进入代理流。
当事件被确认时,引用计数减少,当 它达到零,消费者将认为该偏移量“完成”并且 可以提交。
下面是过滤器https://faust.readthedocs.io/en/latest/userguide/streams.html#id13:
@app.agent() async def process(stream): async for value in stream.filter(lambda: v > 1000).group_by(...): ...
我使用了一些复杂的过滤器,但结果将流分成两部分,用于两个具有完全不同逻辑的代理。 (我不使用 group_by)
如果两个代理一起工作,一切正常。但是,如果我停止它们并重新启动它们,它们将从头开始处理流。因为每一个事件都没有得到代理人之一的承认。 如果我确认每个代理中的所有事件,而如果其中一个代理不会启动,那么第二个代理将清除该主题。 (如果一个被粉碎并重新启动,指挥将看到三个订阅者,因为它正在等待粉碎的代理响应 20 分钟。
我只想将事件分为两部分。在这种情况下如何进行适当的同步?
【问题讨论】:
-
两者都使用相同的组ID?你为消费者配置的起始偏移量是多少?
-
@cricket_007 我刚刚从默认文档的示例开始。
-
好吧,好吧。使用 group by 不是过滤器,您究竟需要同步什么?
-
@cricket_007 我使用
filter,但我不使用group_by,因为我没有一个字段来拆分流。过滤器功能在三个字段上具有条件,以将流拆分为两个代理的两个不相交部分。 -
好的,问题出在哪里?
标签: python apache-kafka faust