【问题标题】:Two agents with different filters on one kafka topic. Acknowledgment in Faust Stream两个代理在一个 kafka 主题上具有不同的过滤器。浮士德流中的确认
【发布时间】: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


【解决方案1】:

faust 过滤在确认过滤出的事件时存在一些错误。我建议在从流中消费时不要使用fault.filter() 功能,而是使用简单的if...then...else 语句样式,类似于以下内容:

@app.agent(topic)
async def process(stream):
    async for event in stream:
        if event.amount >= 300.0:
            yield event

【讨论】:

    猜你喜欢
    • 2019-11-08
    • 2021-06-19
    • 2020-01-25
    • 2022-12-03
    • 1970-01-01
    • 2018-12-10
    • 2015-10-22
    • 2021-05-05
    • 1970-01-01
    相关资源
    最近更新 更多