【问题标题】:SI kafka-Message flow based on router conditionSI kafka-基于路由器条件的消息流
【发布时间】:2015-09-20 16:05:48
【问题描述】:

要求:期望接收/处理inMessageHandler中的inMessage和emailMessageHandler中的emailMessage。

问题:(如果消费者配置的消费者组 ID 不同)消息在这两种情况下都流向消费者服务处理程序的转换方法,但在 emailMessageHandler 中的 emailMessage 的情况下没有得到它。(但在 inMessageHandler 中出现 inMessage )。

问题:(如果消费者配置的消费者组 ID 相同)消息从 ConsumerServiceHandler 的转换方法流向 emailMessageHandler 中的 emailMessage,但根本没有收到 inMessageHandler 中的 inMessage(甚至在 ConsumerServiceHandler 的转换方法中也没有)。

您能否建议这里有什么问题。如何根据主题 id 接收不同服务类的不同消息进行处理?

请在下面找到配置。

    <int-kafka:producer-context id="kafkaProducerContext">
    <int-kafka:producer-configurations>
        <int-kafka:producer-configuration broker-list="localhost:9092"
                   key-class-type="java.lang.String"
                   value-class-type="com.test.EmailMessageVo"
                   topic="emailMessag_topic"
                   value-encoder="emailvalueEncoder"
                   key-encoder="kafkaSerializer"
                   compression-type="none"/>
        <int-kafka:producer-configuration broker-list="localhost:9092"
           key-class-type="java.lang.String"
               value-class-type="com.test.InMessageVo"
                   topic="inMessage_topic"
           value-encoder="invalueEncoder"
                   key-encoder="kafkaSerializer"
                   compression-type="none"/>
    </int-kafka:producer-configurations>
</int-kafka:producer-context>



 <bean id="invalueEncoder" class="org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder">
   <constructor-arg value="com.test.InMessageVo" />

  <bean id="emailvalueEncoder" class="org.springframework.integration.kafka.serializer.avro.AvroReflectDatumBackedKafkaEncoder">
   <constructor-arg value="com.test.EmailMessageVo" />

   <int-kafka:inbound-channel-adapter id="kafkaInboundChannelAdapter"
        kafka-consumer-context-ref="consumerContext"
        auto-startup="false"
        channel="inputFromKafka">
    <int:poller fixed-delay="1" time-unit="MILLISECONDS"/>
</int-kafka:inbound-channel-adapter>



 <int:channel id="inputFromKafka">
   <int:queue />
   </int:channel>


<int:channel id="receiveMessageFromKafka">
  <int:queue />
   </int:channel>



    <int-kafka:consumer-context id="consumerContext"
        consumer-timeout="1000"
        zookeeper-connect="zookeeperConnect" consumer-properties="consumerProperties">
    <int-kafka:consumer-configurations>
        <int-kafka:consumer-configuration group-id="default1"
                value-decoder="emailvalueDecoder"
                key-decoder="kafkaReflectionDecoder"
                max-messages="5000">
            <int-kafka:topic id="test1" streams="4"/>
        </int-kafka:consumer-configuration>

        <int-kafka:consumer-configuration group-id="default2"
                value-decoder="invalueDecoder"
                key-decoder="kafkaReflectionDecoder"
                max-messages="50">
        <int-kafka:topic id="test2" streams="4"/>
        </int-kafka:consumer-configuration>
    </int-kafka:consumer-configurations>
</int-kafka:consumer-context>


  <bean id="emailvalueDecoder" class="org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaDecoder">
    <constructor-arg value="com.test.EmailMessageVo" />
</bean>


  <bean id="invalueDecoder" class="org.springframework.integration.kafka.serializer.avro.AvroSpecificDatumBackedKafkaDecoder">
    <constructor-arg value="com.test.InMessageVo" />
</bean>
   <bean id="consumerService" class="com.test.ConsumerServiceHandler" />

   <int:service-activator id="advicedSa" input-channel="inputFromKafka" ref="consumerService"method="testConsumerCircuitBreaker" outputchannel="receiveMessageFromKafka">

   <int:channel id="inMessage_topic_channel">
   <int:queue />

   <int:channel id="emailMessage_topic_channel">
    <int:queue />

    <bean id="inMessageService" class="com.test.inMessageHandler" />

     <bean id="emailMessageService" class="com.test.emailMessageHandler" />

      <int:service-activator id="advicedSa1" input-channel="inMessage_topic_channel" ref="inMessageService"method="execute">

      <int:service-activator id="advicedSa2" input-channel="emailMessage_topic_channel" ref="emailMessageService"method="execute" >

       <int:router input-channel="receiveMessageFromKafka" expression="headers.topic">
       <int:mapping value="inMessage_topic" channel="inMessage_topic_channel"/>
       <int:mapping value="emailMessage_topic" channel="emailMessage_topic_channel"/>

【问题讨论】:

    标签: spring-integration apache-kafka


    【解决方案1】:

    您应该与我们分享您的ConsumerServiceHandler 所做的事情。仅仅因为您的逻辑有点奇怪,将&lt;int-kafka:inbound-channel-adapter&gt; 的结果作为Message&lt;Map&lt;String, Map&lt;Integer, List&lt;Object&gt;&gt;&gt;&gt;,它是Message&lt;?&gt; with thepayload` 作为分区部分的主题及其消息的映射

    使用相同的group-id,我们(你?)最终会遇到最后一个配置胜过前一个配置的问题:

    consumerConfigurationsMap.put(consumerConfiguration.getAttribute("group-id"),
                    consumerConfigurationBeanDefinition);
    

    【讨论】:

    • 感谢@ArtemBilan 的友好回复,请在此处找到consumerServiceHandler:gist.github.com/anonymous/d8cf04dcabd8e4463d33 ConsumerServiceHandler:删除主题/分区冗余细节/设置主题。发送的消息:Message.withPayload(inMessageVO)/Message.withPayload(emailMessageVO),期望相同但格式如下:Message>>> 总体意图:路由不同的消息( inMessageVO/emailMessageVO) 到 diff 处理程序进行处理请建议,如果它可以以更好的方式完成。
    猜你喜欢
    • 1970-01-01
    • 2016-06-26
    • 2021-10-03
    • 1970-01-01
    • 1970-01-01
    • 2016-05-04
    • 2020-01-12
    • 2019-05-12
    • 1970-01-01
    相关资源
    最近更新 更多