Потребитель конфигурации spring-integration-kafka для получения сообщения из указанного раздела

Я начал использовать spring-integration-kafka в своем проекте, и я могу создавать и использовать сообщения из Kafka. Но теперь я хочу создать сообщение для определенного раздела, а также использовать сообщение из определенного раздела.

Пример: я хочу создать сообщение для раздела 3, и потребитель будет получать сообщение только из раздела 3.

До сих пор в моей теме было 8 разделов, и я могу создавать сообщения для определенного раздела, но я еще не нашел способ настроить, чтобы потребитель получал сообщение только из определенного раздела.

Поэтому любые предложения о том, как мне настроить потребителя с помощью spring-integration-kafka или что-то еще, что нужно сделать с классом KafkaConsumer.java, чтобы он мог получать сообщение из определенного раздела.

Спасибо.

Вот мой код:

кафка-производитель-context.xml

<int:publish-subscribe-channel id="inputToKafka" />

<int-kafka:outbound-channel-adapter
    id="kafkaOutboundChannelAdapter" kafka-producer-context-ref="kafkaProducerContext"
    auto-startup="true" order="1" channel="inputToKafka" />
<int-kafka:producer-context id="kafkaProducerContext"
    producer-properties="producerProps">
    <int-kafka:producer-configurations>
        <int-kafka:producer-configuration 
            broker-list="127.0.0.1:9092,127.0.0.1:9093,127.0.0.1:9094"
            async="true" topic="testTopic"
            key-class-type="java.lang.String" 
            key-encoder="encoder"
            value-class-type="java.lang.String" 
            value-encoder="encoder"
            partitioner="partitioner"
            compression-codec="default" />
    </int-kafka:producer-configurations>
</int-kafka:producer-context>

<util:properties id="producerProps">
    <prop key="queue.buffering.max.ms">500</prop>
    <prop key="topic.metadata.refresh.interval.ms">3600000</prop>
    <prop key="queue.buffering.max.messages">10000</prop>
    <prop key="retry.backoff.ms">100</prop>
    <prop key="message.send.max.retries">2</prop>
    <prop key="send.buffer.bytes">5242880</prop>
    <prop key="socket.request.max.bytes">104857600</prop>
    <prop key="socket.receive.buffer.bytes">1048576</prop>
    <prop key="socket.send.buffer.bytes">1048576</prop>
    <prop key="request.required.acks">1</prop>
</util:properties>

<bean id="encoder"
    class="org.springframework.integration.kafka.serializer.common.StringEncoder" />

<bean id="partitioner" class="org.springframework.integration.kafka.support.DefaultPartitioner"/>

<task:executor id="taskExecutor" pool-size="5"
    keep-alive="120" queue-capacity="500" />

KafkaProducer.java

public class KafkaProducer {

private static final Logger logger = LoggerFactory
        .getLogger(KafkaProducer.class);

@Autowired
private MessageChannel inputToKafka;

public void sendMessage(String message) {

    try {
        inputToKafka.send(MessageBuilder.withPayload(message)
                    .setHeader(KafkaHeaders.TOPIC, "testTopic")
                    .setHeader(KafkaHeaders.PARTITION_ID, 3).build());
    } catch (Exception e) {
        logger.error(String.format(
                "Failed to send [ %s ] to topic %s ", message, topic),
                e);
    }
}

}

файл kafka-consumer-context.xml

<int:channel id="inputFromKafka">
    <int:dispatcher task-executor="kafkaMessageExecutor" />
</int:channel>

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

<int-kafka:inbound-channel-adapter
    id="kafkaInboundChannelAdapter" 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>


<bean id="consumerProperties"
    class="org.springframework.beans.factory.config.PropertiesFactoryBean">
    <property name="properties">
        <props>
            <prop key="auto.offset.reset">smallest</prop>
            <prop key="socket.receive.buffer.bytes">1048576</prop>
            <prop key="fetch.message.max.bytes">5242880</prop>
            <prop key="auto.commit.interval.ms">1000</prop>
        </props>
    </property>
</bean>

<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="defaultGrp" max-messages="20000">
            <int-kafka:topic id="testTopic" streams="3" />
        </int-kafka:consumer-configuration>
    </int-kafka:consumer-configurations>
</int-kafka:consumer-context>

<task:executor id="kafkaMessageExecutor" pool-size="0-10"
    keep-alive="120" queue-capacity="500" />

<int:outbound-channel-adapter channel="inputFromKafka"
    ref="kafkaConsumer" method="processMessage" />

KafkaConsumer.java

public class KafkaConsumer {

private static final Logger log = LoggerFactory
        .getLogger(KafkaConsumer.class);

@Autowired
KafkaService kafkaService;

public void processMessage(Map<String, Map<Integer, List<byte[]>>> msgs) {
    for (Map.Entry<String, Map<Integer, List<byte[]>>> entry : msgs
            .entrySet()) {
        log.debug("Topic:" + entry.getKey());
        ConcurrentHashMap<Integer, List<byte[]>> messages = (ConcurrentHashMap<Integer, List<byte[]>>) entry
                .getValue();
        log.debug("\n**** Partition: \n");
        Set<Integer> keys   = messages.keySet();
        for (Integer i : keys)
            log.debug("p:"+i);
        log.debug("\n**************\n");
        Collection<List<byte[]>> values = messages.values();
        for (Iterator<List<byte[]>> iterator = values.iterator(); iterator
                .hasNext();) {
            List<byte[]> list = iterator.next();
            for (byte[] object : list) {
                String message = new String(object);
                log.debug("Message: " + message);
                try {
                    kafkaService.receiveMessage(message);
                } catch (Exception e) {
                    log.error(String.format("Failed to process message %s",
                            message));
                }
            }
        }

    }
}
}

Итак, моя проблема здесь. Когда я отправляю сообщение в раздел 3 или любой раздел, KafkaConsumer всегда получает сообщение. Все, что я хочу: KafkaConsumer будет получать сообщение только из раздела 3, а не из другого раздела.

Спасибо еще раз.


person Sebastien Le    schedule 29.06.2015    source источник
comment
Я столкнулся с такой же проблемой, как ты. вы нашли какое-нибудь решение?   -  person user3359139    schedule 12.09.2015
comment
Я также хочу интегрировать Kafka со Spring, не могли бы вы поделиться любым репо, откуда я могу скачать рабочие коды   -  person Sankalp    schedule 14.10.2016


Ответы (1)


Вам необходимо использовать адаптер канала, управляемый сообщениями < / а>.

Как вариант, KafkaMessageListenerContainer может принимать аргумент массива org.springframework.integration.kafka.core.Partition для указания тем и их пары разделов.

Вам необходимо подключить контейнер слушателя, используя этот конструктор и передайте его адаптеру с помощью атрибута listener-container.

Мы обновим файл readme примером.

person Gary Russell    schedule 29.06.2015
comment
Спасибо, Гэри, я пытаюсь это сделать и жду, например. - person Sebastien Le; 30.06.2015
comment
Привет, Гэри, похоже, что адаптер канала, управляемый сообщениями, использует Kafka SimpleConsumer. Итак, мой вопрос: могу ли я настроить потребителя для получения сообщения из указанного раздела в Kafka High-level Consumer? Потому что у меня несколько потребителей. Спасибо. - person Sebastien Le; 30.06.2015
comment
Нет; для выбора конкретных тем вам понадобится адаптер, управляемый сообщениями. - person Gary Russell; 30.06.2015
comment
Привет, Гэри. У меня есть еще одна проблема: когда я пытаюсь отправить около 10 сообщений в секунду, похоже, что потребитель может получить сообщение немедленно. Но когда я пытаюсь использовать большой размер, около 500 месс / 1 с, у потребителя есть время задержки, около 3-5 секунд, чтобы начать прием сообщения. Не знаю почему. У вас есть идеи или предложения, которые помогут мне решить эту проблему? Огромное спасибо. - person Sebastien Le; 14.07.2015