【发布时间】:2020-10-29 07:48:08
【问题描述】:
- 我有两个使用 Spring Boot 用 Java 编写的微服务。
- 我使用 Kafka,通过 Spring Cloud Stream Kafka,在它们之间发送消息。
- 我需要发送一个自定义标头,但直到现在都没有成功。
- 我已阅读并尝试了我在 Internet 和 Spring Cloud Stream 文档上找到的大部分内容...
...我仍然无法让它工作。
这意味着我永远不会在接收器中收到消息,因为无法找到标头并且不能为空。 我怀疑标题从未写在消息中。现在我正在尝试用 Kafkacat 验证这一点。
欢迎任何帮助
提前致谢。
------信息------
这里是发件人代码:
@SendTo("notifications")
public void send(NotificationPayload payload, String eventId) {
var headerMap = Collections.singletonMap("EVENT_ID",
eventId.getBytes(StandardCharsets.UTF_8));
MessageHeaders headers = new MessageHeaders(headerMap);
var message = MessageBuilder.createMessage(payload, headers);
notifications.send(message);
}
其中notifications 是MessageChannel
这里是消息发送者的相关配置。
spring:
cloud:
stream:
defaultBinder: kafka
bindings:
notifications:
binder: kafka
destination: notifications
contentType: application/x-java-object;type=com.types.NotificationPayload
producer:
partitionCount: 1
headerMode: headers
kafka:
binder:
headers: EVENT_ID
我也试过headers: "EVENT_ID"
这是接收部分的代码:
@StreamListener("notifications")
public void receiveNotif(@Header("EVENT_ID") byte[] eventId,
@Payload NotificationPayload payload) {
var eventIdS = new String((byte[]) eventId, StandardCharsets.UTF_8);
...
// do something with the payload
}
以及接收部分的配置:
spring:
cloud:
stream:
kafka:
bindings:
notifications:
consumer:
headerMode: headers
版本
<spring-cloud-stream-dependencies.version>Horsham.SR4</spring-cloud-stream-dependencies.version>
<spring-cloud-stream-binder-kafka.version>3.0.4.RELEASE</spring-cloud-stream-binder-kafka.version>
<spring-cloud-schema-registry.version>1.0.4.RELEASE</spring-cloud-schema-registry.version>
<spring-cloud-stream.version>3.0.4.RELEASE</spring-cloud-stream.version>
【问题讨论】:
标签: spring-boot apache-kafka spring-cloud spring-cloud-stream spring-cloud-stream-binder-kafka