如何在Spring Cloud AWS Kinesis Binder中设置动态分区键?
使用Spring Cloud Stream Supplier为Kinesis动态设置分区键
问题
我尝试使用Supplier/Consumer组件从Kinesis数据流生产和消费消息,请问是否可实现动态添加分区键?
我的代码实现
Java代码:
private BlockingQueue<Message> messages = new LinkedBlockingQueue<>(); @Bean public Supplier<Message<String>> produceMessages() { return () -> this.messages.poll(); } @Override public void produce(main.Test request, StreamObserver<Test> response) { Message input = MessageBuilder.withPayload(request.getMessage()) .setHeader("partitionKey", "los").build(); this.messages.offer(input); response.onCompleted(); }
application.properties配置:
spring.cloud.stream.bindings.produceMessages-out-0.producer.partitionKeyExpression=headers['partitionKey']
回答
可以实现动态添加分区键,你的代码写法本身就是正确的实现方案。
- 你通过
MessageBuilder构建消息时,动态为每个消息设置partitionKey消息头,这一步实现了分区键的动态赋值; - 配置项
partitionKeyExpression=headers['partitionKey']指定Kinesis生产者从消息头中读取分区键的值,框架会自动将该值作为Kinesis记录的分区键发送到数据流中。
这种方式完全满足动态为不同消息设置不同分区键的需求,不需要额外修改。
内容的提问来源于stack exchange,提问作者Kanika
相关产品推荐
相关产品推荐

