【问题标题】:Spring integration kafka inbound adapter aggregating incoming messagesSpring集成kafka入站适配器聚合传入消息
【发布时间】:2015-08-11 02:44:51
【问题描述】:

我有一个int-kafka:outbound-channel-adapter,我用它向kafka 发送消息,然后使用int-kafka:inbound-channel-adapter 接收消息。通信似乎工作正常,我能够发送和接收消息,但格式有点奇怪。我将单独的消息单独发送到我的出站适配器,但是当我收到消息时,我会收到一条消息,其中所有消息都聚合到该消息的有效负载中。

这是我收到消息时消息负载的样子

[payload={mytopic={0=[字符串消息 1, 字符串消息 2, 字符串消息 3, 字符串消息 4, 字符串消息 5, .........]}}, headers= {id=3934de02-1f42-ab90-6aa5-9c15f3cd0b6e,时间戳=1439260669762}]

接收集成流程如下所示

<int-kafka:inbound-channel-adapter
    id="kafkaInboundAdapter" kafka-consumer-context-ref="consumerContext"
    auto-startup="true" channel="inputFromKafka">
    <int:poller fixed-delay="10" time-unit="MILLISECONDS"
        max-messages-per-poll="5" />
</int-kafka:inbound-channel-adapter>

<int:channel id="inputFromKafka" />

<int:service-activator id="kakfaMessageHandler"
    input-channel="inputFromKafka">
    <bean class="com...broker.MessageHandler"></bean>
</int:service-activator>

我收到所有消息汇总在一条 spring 集成消息中而不是发送到 kafka 时的单独消息的任何原因。

【问题讨论】:

    标签: spring-integration apache-kafka


    【解决方案1】:

    KafkaHighLevelConsumerMessageSource 与许多其他轮询MessageSource&lt;?&gt; 一样设计:通过一次轮询获取数据并将其作为List&lt;?&gt; 返回。

    在这种情况下,我们从 Kafka stream 得到这个结果:

    Message<Map<String, Map<Integer, List<Object>>>>
    

    payload 是 Kafka 的 Map topics 和 partitions 和 messages 的映射。

    如果您在consumerContext 上仅使用一个topic,您可以简单地将顶级Map 转换为其partitions 映射。或者,如果您只有一个 partition ,甚至可以继续转换到 payloads 列表。最后你可以得到splitter

    如果您想尽快收到来自该主题的消息,就像它们出现在那里一样快,您应该查看&lt;int-kafka:message-driven-channel-adapter&gt;

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-08-05
      • 2018-04-08
      • 2015-03-21
      • 2015-07-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多