使用PostgreSQL实现Kinesis Checkpointing与锁时重复消费问题求助
问题:Spring Cloud Stream Kinesis绑定器用PostgreSQL做Checkpoint后消息重复消费
环境与依赖
- 使用依赖:
org.springframework.cloud:spring-cloud-stream-binder-kinesis:4.0.2 - 部署架构:两个负载均衡的Pod,连接同一个Kinesis流与PostgreSQL数据库
- 需求:用现有PostgreSQL替代DynamoDB实现Checkpointing和锁功能,降低资源成本
application.yml配置
spring: cloud: aws: credentials: sts: web-identity-token-file: <Where i had given the token file path> role-arn: <Where i had given the assume role arn> role-session-name: RoleSessionName region: static: <where i had given my aws region> dualstack-enabled: false stream: kinesis: binder: auto-create-stream: false min-shard-count: 1 bindings: input-in-0: destination: test-test.tst.v1 content-type: text/json
Kinesis消费处理代码
@Configuration public class KinesisConsumerBinder{ @Bean public Consumer<Message<String>> input(){ return message ->{ System.out.println("Data from Kinesis:"+message.getPayload()); //Process the message got from Kinesis } } }
已实现方案与问题
已通过自定义实现将PostgreSQL作为Checkpoint和锁的存储介质,但运行时出现每条Kinesis消息被重复插入业务表的问题。
以下是Checkpoint和锁表的示例数据:
Checkpoint表(int_metadata_store)
| metadata_key | metadata_value | region |
|---|---|---|
| anonymous.6168ad-9a13-5b75-96b9-996f340dfd:test-test.tst.v1:shardId-000000001 | 8954146227109765333558818710934653698471926607 | DEFAULT |
| anonymous.7168ad-8b13-5c75-96n9-996f340dfd:test-test.tst.v1:shardId-000000001 | 8954146227109765333558818710934653698471926607 | DEFAULT |
锁表(int_lock)
| lock_key | region | client_id | created_date |
|---|---|---|---|
| 44444444-7df8-2222-db33-c0bbb07bbbb8 | DEFAULT | a9a9a9a9-8888-4444-777d-99a33aa20122 | 2024-03-26 09:28:29.662 |
| 55555555-dddd-3333-b000-8810ae858a84 | DEFAULT | a9a9a9a9-8888-4444-777d-99a33aa20122 | 2024-03-26 09:28:29.662 |
| 66666666-eeee-4444-a805-6138a0c59976 | DEFAULT | a9a9a9a9-8888-4444-777d-99a33aa20122 | 2024-03-26 09:28:29.662 |
期望
每条Kinesis消息仅在业务表中保存一次,需解决重复消费问题。
内容的提问来源于stack exchange,提问作者Jerrin Paul
相关产品推荐
相关产品推荐

