代码之家  ›  专栏  ›  技术社区  ›  maverick

消息预处理(主题-主题)-Kafka Connect API vs.Streams vs Kafka Consumer?

  •  0
  • maverick  · 技术社区  · 8 年前

    我们需要对从一个主题到另一个主题的每条消息进行一些预处理(使用不同的密钥解密/重新加密)。

    我一直在研究使用Kafka Connect,因为它提供了很多现成的好东西(配置管理、偏移存储、错误处理等)。

    但我也觉得 SourceConnector SinkConnector 只是在两个主题之间移动数据,而这两个接口都不能 Topic A -> (Connector) -> Topic B .这是正确的方法吗?我应该使用 下沉接头 独自一人 SourceTask.put() 所有的逻辑都写给卡夫卡了吗?

    其他选项包括 KafkaConsumer/Producer and或Streams,但它们将需要自己的实例来运行逻辑,而不是偏移重试错误处理。

    1 回复  |  直到 8 年前
        1
  •  2
  •   OneCricketeer Gabriele Mariotti    8 年前

    提供了许多现成的好东西(配置管理、偏移存储、错误处理等)

    配置管理不应该比重新部署应用程序更难,但这取决于您可能拥有或没有的任何版本控制或CI/CD管道。

    卡夫卡制作者/消费者和流提供偏移管理,您只需将其配置为执行除默认值以外的任何操作。

    错误处理有很好的文档记录,如果您关心错误的检测,请不要忘记。连接本身将在严重错误情况下停止使用和生成消息,而不会重试或跳过消息。

    这两个接口都不能 Topic A -> (Connector) -> Topic B “”

    你见过Confluent Replicator(许可产品)吗?那个 卡夫卡连接了两个主题。

    否则,你见过MirrorMaker吗?这是一个生产者-消费者对,通常用于在各个集群之间复制数据,但可以与相同的源和目标设置一起使用。您只需要确保您没有创建反馈循环。您需要对其应用“自定义逻辑”(并更改主题名称),其中 called a Handler class that is placed on your Kafka classpath

    bin/kafka-mirror-maker.sh
    
    ...
    
    --message.handler <String: A custom      Message handler which will process
      message handler of type                  every record in-between consumer and
      MirrorMakerMessageHandler>               producer.
    --message.handler.args <String:          Arguments used by custom message
      Arguments passed to message handler      handler for mirror maker.
      constructor.>
    

    Confluent MirrorMaker documentation
    Kafka MirrorMaker documentation


    没有什么可以阻止您实现Connect API,它可能比没有外部集群管理器的Kafka Streams应用程序更易于管理。此外,由于Connect是一个Java库,所以理论上您可以在其中内部使用Streams库。

    推荐文章