【发布时间】:2015-10-06 07:24:53
【问题描述】:
下面是我的配置
<int-kafka:inbound-channel-adapter id="kafkaInboundChannelAdapter"
kafka-consumer-context-ref="consumerContext"
auto-startup="true"
channel="inputFromKafka">
<int:poller fixed-delay="1" time-unit="MILLISECONDS" />
</int-kafka:inbound-channel-adapter>
inputFromKafka下面经过改造
public Message<?> transform(final Message<?> message) {
System.out.println( "KAFKA Message Headers " + message.getHeaders());
final Map<String, Map<Integer, List<Object>>> origData = (Map<String, Map<Integer, List<Object>>>) message.getPayload();
// some code to figure-out the nonPartitionedData
return MessageBuilder.withPayload(nonPartitionedData).build();
}
上面的打印语句只打印两个一致的标题
KAFKA Message Headers {id=9c8f09e6-4b28-5aa1-c74c-ebfa53c01ae4, timestamp=1437066957272}
在发送 Kafka 消息时,传递了一些标头,包括 KafkaHeaders.MESSAGE_KEY,但我也没有回复,想知道是否有办法完成此操作?
【问题讨论】:
标签: java spring spring-integration apache-kafka