【发布时间】:2019-03-15 13:50:42
【问题描述】:
我正在尝试使用 Spring Cloud Stream 创建一个 kafka Consumer,以便侦听一个在任何 Spring 上下文之外构建的带有自定义标头 (operationType) 的 Kafka 消息。
我正在使用 Spring Boot 1.5.x / Spring Cloud Egdware.SR5 和 1.1.1 版本的 kafka-client 和 kafka_2.11。
我的 Listener 类包含此方法
@StreamListener(value = "dataset-changed", condition = "headers['operationType']=='UPDATE'")
public void onEvent(@Payload DatasetChangedMessage payload) {
// my code should be execute only if the header operationType == UPDATE
}
Spring Cloud Stream 配置是
spring.cloud.stream:
bindings:
dataset-changed:
group: preparation
content-type: application/json
destination: dataset-changed
consumer:
headerMode: raw
configuration:
key.deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer
value.deserializer: org.apache.kafka.common.serialization.StringDeserializer
producer 是一个带有 kafka-client:1.1.1 库的简单 java doc
Properties producerConfig = new Properties();
producerConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
producerConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<byte[], String> producer = new KafkaProducer<>(producerConfig);
// Headers (to condition the kafka listener)
final List<Header> headers = new ArrayList<>();
headers.add(new RecordHeader("operationType", "UPDATE".getBytes()));
ProducerRecord<byte[], String> record =
new ProducerRecord<>("dataset-changed", 0, "111".getBytes(), getJsonPayload(), headers);
Future<RecordMetadata> future = producer.send(record);
future.get();
producer.close();
当我生成 kafka 消息时,我有这种警告
2019-03-15 14:48:32.103 WARN [tdp-preparation,1e24b9764ef9bb14,1e24b9764ef9bb14,false] 34760 --- [ -L-1] .DispatchingStreamListenerMessageHandler : Cannot find a @StreamListener matching for message with id: ea27a446-69da-7b8d-1b94-50b46a40dfde
而 operationType 标头存在
【问题讨论】:
标签: java spring-cloud-stream spring-kafka