【问题标题】:Java DSL Equivalent of Spring Integration Kafka Endpoint Configured in XML以 XML 配置的 Spring 集成 Kafka 端点的 Java DSL 等效项
【发布时间】:2016-03-30 05:11:26
【问题描述】:

我有以下用于 Kafka 出站通道适配器的 XML 配置:

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

</int-kafka:outbound-channel-adapter>
<task:executor id="taskExecutor"
               pool-size="5-25"
               queue-capacity="20"
               keep-alive="120"/>

这很好用。我试图在 Java DSL 中复制它,但我不能走得太远。到目前为止,我只有这个:

.handle(Kafka.outboundChannelAdapter(kafkaConfig)
        .addProducer(producerMetadata, brokerAddress)
        .get());

我不知道如何在 DSL 中添加 taskExecutorpoller

感谢任何关于如何将这些融入我的整体IntegrationFlow 的见解。

【问题讨论】:

    标签: java spring spring-integration apache-kafka


    【解决方案1】:

    Spring Integration 组件(例如 &lt;int-kafka:outbound-channel-adapter&gt;)由两个 bean 组成:AbstractEndpoint 用于接受来自 input-channel 的消息,MessageHandler 用于处理消息。

    所以,Kafka.outboundChannelAdapter() 大约是 MessageHandler。任何其他特定于端点的属性都取决于.handle() EIP 方法的第二个Consumer&lt;GenericEndpointSpec&lt;H&gt;&gt; endpointConfigurer 参数:

    .handle(Kafka.outboundChannelAdapter(kafkaConfig)
        .addProducer(producerMetadata, brokerAddress),
               e -> e.id("kafkaOutboundChannelAdapter")
                     .poller(p -> p.fixedDelay(1000, TimeUnit.MILLISECONDS)
                                            .receiveTimeout(0)
                                            .taskExecutor(this.taskExecutor)));
    

    更多信息请参见Reference Manual

    【讨论】:

      猜你喜欢
      • 2021-08-02
      • 1970-01-01
      • 2018-06-21
      • 1970-01-01
      • 1970-01-01
      • 2017-06-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多