【问题标题】:kafka Message driven adapter not workingkafka消息驱动适配器不工作
【发布时间】:2015-11-05 19:11:14
【问题描述】:

想使用 kafka:message-driven-channel-adapter 来消费消息。

在以下渠道生成消息: headers['topic'] = inMsge_topic,emailMesge_topic

在互联网上也没有找到同样的好例子。请提出建议。

使用时工作正常
int-kafka:inbound-channel-adapter 但它需要轮询。(不想轮询)

下面是使用的配置:

 <int-kafka:outbound-channel-adapter id="kafkaOutboundChannelAdapter" kafka-producer-context-ref="kafkaProducerContext"
                                                                 auto-startup="true" channel="inputToKafka" message-key="kafka_messageKey">
          <int:poller fixed-delay="1000" time-unit="MILLISECONDS" receive-timeout="0" task-executor="taskExecutor" /> 
   </int-kafka:outbound-channel-adapter>

   <int-kafka:producer-context id="kafkaProducerContext"
          producer-properties="producerProperties">
          <int-kafka:producer-configurations>
                <int-kafka:producer-configuration
                       broker-list="${kafka.producer.brokerList}" topic="headers['topic']" key-class-type="java.lang.String"
                       value-class-type="com.vo.MessageVO"
                       value-encoder="kafkaEncoder" key-encoder="kafkaKeyEncoder"
                       compression-type="none" />
          </int-kafka:producer-configurations>
   </int-kafka:producer-context>    
   <int:channel id="inputToKafka">
     <int:queue />
   </int:channel>

  <int:channel id="inputFromKafka">
   </int:channel>
 <bean id="kafkaConfiguration" class="org.springframework.integration.kafka.core.ZookeeperConfiguration">
   <constructor-arg ref="zookeeperConnect"/>

 <bean id="connectionFactory" class="org.springframework.integration.kafka.core.DefaultConnectionFactory">
    <constructor-arg ref="kafkaConfiguration"/>

      <int-kafka:message-driven-channel-adapter
        id="adapter"
        channel="inputFromKafka"
        connection-factory="connectionFactory"
        key-decoder="kafkaKeyDecoder"
        payload-decoder="kafkaDecoder"          
        max-fetch="100"
        topics="inMsge_topic"/>

        <int-kafka:message-driven-channel-adapter
        id="adapter1"
        channel="inputFromKafka"
        connection-factory="connectionFactory"
        key-decoder="kafkaKeyDecoder"
        payload-decoder="kafkaDecoder"          
        max-fetch="100"
        topics="emailMesge_topic"/>


    <int-kafka:zookeeper-connect id="zookeeperConnect"
    zk-connect="localhost:2181" zk-connection-timeout="6000"
    zk-session-timeout="400" zk-sync-time="200" />

日志输出:

20:01:13.130 [pool-5-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取[topic='inMsge_topic', id=2]@0 20:01:13.131 [pool-5-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=4]@0 20:01:13.131 [pool-5-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=1]@0 20:01:13.131 [pool-5-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=0]@0 20:01:13.131 [pool-5-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=3]@1913 20:01:13.134 [pool-11-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=1]@0 20:01:13.134 [pool-11-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=3]@1913 20:01:13.134 [pool-11-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=4]@0 20:01:13.134 [pool-11-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=2]@0 20:01:13.134 [pool-11-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='inMsge_topic', id=0]@0 20:01:13.158 [pool-7-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=1]@0 20:01:13.158 [pool-7-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=3]@334 20:01:13.158 [pool-7-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=0]@0 20:01:13.158 [pool-7-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=4]@0 20:01:13.158 [pool-7-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=2]@0 20:01:13.164 [pool-13-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=1]@0 20:01:13.164 [pool-13-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=4]@0 20:01:13.164 [pool-13-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=0]@0 20:01:13.164 [pool-13-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=3]@334 20:01:13.164 [pool-13-thread-1] 调试 o.s.i.kafka.core.DefaultConnection - 从分区读取 [topic='emailMesge_topic', id=2]@0

【问题讨论】:

    标签: spring spring-integration apache-kafka


    【解决方案1】:

    看看Spring Integration kafka sample - 它使用 Java 配置而不是 XML,但它同时显示出站和消息驱动的适配器。

    【讨论】:

    • 嗨,加里,感谢您的建议.. kafkaMessageDrivenChannelAdapter.setOutputChannel(received()); @Bean public PollableChannel received() { return new QueueChannel(); } 它会轮询消息吗?寻找监听器配置。
    • 我们也可以使用 PublishSubscribeChannel/DirectChannel 等。是吗?
    • 是的。你可以。 `kafkaMessageDrivenChannelAdapter.setOutputChannel(received());` 表示从 Kafka 主题检索消息后将发送到哪里。
    • 嗨@Artem 感谢您的建议。尝试使用上面提到的 Java 配置示例/链接,此处出现错误提示。 (gist.github.com/anonymous/75d51a3d1fbe73e0b0ee)。请提出建议。
    • 看起来您应该提出一个新问题并就此问题分享您的 Web 应用程序配置。您可以只使用不是由根 webapp 上下文启动的单独上下文中的consumer
    猜你喜欢
    • 2017-12-23
    • 2017-04-29
    • 1970-01-01
    • 2021-03-20
    • 1970-01-01
    • 2014-01-23
    • 2011-04-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多