如何用PostgreSQL替代DynamoDB实现Kinesis Binder的Checkpoint与锁?
用PostgreSQL替代DynamoDB实现Spring Cloud Stream Kinesis的Checkpoint与分布式锁
基于spring-cloud-stream-binder-kinesis:4.0.2,你可以通过实现框架提供的CheckpointStore和LockRegistry接口,直接复用现有PostgreSQL实现Checkpoint存储和分布式锁,无需手动处理消息头和Checkpoint逻辑。
步骤1:创建PostgreSQL数据库表
首先在你的PostgreSQL中创建Checkpoint存储表和分布式锁表:
-- Checkpoint存储表:关联流、分片、消费组和最新处理的序列号 CREATE TABLE kinesis_checkpoint ( stream_name VARCHAR(255) NOT NULL, shard_id VARCHAR(255) NOT NULL, consumer_group VARCHAR(255) NOT NULL, sequence_number VARCHAR(255) NOT NULL, last_updated TIMESTAMP DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (stream_name, shard_id, consumer_group) ); -- 分布式锁表:用于分片的独占消费控制 CREATE TABLE kinesis_lock ( lock_key VARCHAR(255) PRIMARY KEY, owner VARCHAR(255) NOT NULL, expiration_time TIMESTAMP NOT NULL );
步骤2:实现Checkpoint存储逻辑
2.1 定义Checkpoint实体类与Repository
import jakarta.persistence.Column; import jakarta.persistence.Entity; import jakarta.persistence.Id; import jakarta.persistence.Table; import java.time.LocalDateTime; @Entity @Table(name = "kinesis_checkpoint") public class KinesisCheckpointEntity { @Id @Column(name = "stream_name") private String streamName; @Id @Column(name = "shard_id") private String shardId; @Id @Column(name = "consumer_group") private String consumerGroup; @Column(name = "sequence_number") private String sequenceNumber; @Column(name = "last_updated") private LocalDateTime lastUpdated; // 无参构造器(JPA要求) public KinesisCheckpointEntity() {} // 带参构造器 public KinesisCheckpointEntity(String streamName, String shardId, String consumerGroup) { this.streamName = streamName; this.shardId = shardId; this.consumerGroup = consumerGroup; this.lastUpdated = LocalDateTime.now(); } // Getter和Setter省略 }
import org.springframework.data.jpa.repository.JpaRepository; import java.util.Optional; public interface KinesisCheckpointRepository extends JpaRepository<KinesisCheckpointEntity, String> { Optional<KinesisCheckpointEntity> findByStreamNameAndShardIdAndConsumerGroup(String streamName, String shardId, String consumerGroup); void deleteByStreamNameAndShardIdAndConsumerGroup(String streamName, String shardId, String consumerGroup); }
2.2 实现CheckpointStore接口
该接口由Spring Cloud Stream Kinesis Binder调用,自动处理Checkpoint的读取、保存和删除:
import org.springframework.cloud.stream.binder.kinesis.CheckpointStore; import org.springframework.stereotype.Component; import java.util.Optional; @Component public class PostgresCheckpointStore implements CheckpointStore { private final KinesisCheckpointRepository checkpointRepository; public PostgresCheckpointStore(KinesisCheckpointRepository checkpointRepository) { this.checkpointRepository = checkpointRepository; } @Override public Optional<String> getCheckpoint(String stream, String shard, String consumerGroup) { return checkpointRepository.findByStreamNameAndShardIdAndConsumerGroup(stream, shard, consumerGroup) .map(KinesisCheckpointEntity::getSequenceNumber); } @Override public void saveCheckpoint(String stream, String shard, String consumerGroup, String sequenceNumber) { KinesisCheckpointEntity entity = checkpointRepository.findByStreamNameAndShardIdAndConsumerGroup(stream, shard, consumerGroup) .orElse(new KinesisCheckpointEntity(stream, shard, consumerGroup)); entity.setSequenceNumber(sequenceNumber); entity.setLastUpdated(java.time.LocalDateTime.now()); checkpointRepository.save(entity); } @Override public void deleteCheckpoint(String stream, String shard, String consumerGroup) { checkpointRepository.deleteByStreamNameAndShardIdAndConsumerGroup(stream, shard, consumerGroup); } }
步骤3:实现分布式锁逻辑
多实例部署时,需要用分布式锁保证每个Kinesis分片仅被一个消费者处理,这里基于PostgreSQL的UPSERT实现:
import org.springframework.integration.support.locks.LockRegistry; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.stereotype.Component; import java.time.LocalDateTime; import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; @Component public class PostgresLockRegistry implements LockRegistry { private final JdbcTemplate jdbcTemplate; private final String lockOwner = UUID.randomUUID().toString(); private final java.time.Duration lockTimeout = java.time.Duration.ofMinutes(5); // 锁超时,防止死锁 public PostgresLockRegistry(JdbcTemplate jdbcTemplate) { this.jdbcTemplate = jdbcTemplate; } @Override public Lock obtain(Object lockKey) { return new PostgresDistributedLock(lockKey.toString()); } private class PostgresDistributedLock implements Lock { private final String lockKey; public PostgresDistributedLock(String lockKey) { this.lockKey = lockKey; } @Override public void lock() { while (!tryLock()) { try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IllegalStateException("锁获取被中断", e); } } } @Override public boolean tryLock() { String sql = """ INSERT INTO kinesis_lock (lock_key, owner, expiration_time) VALUES (?, ?, ?) ON CONFLICT (lock_key) DO UPDATE SET owner = EXCLUDED.owner, expiration_time = EXCLUDED.expiration_time WHERE kinesis_lock.expiration_time < NOW() OR kinesis_lock.owner = ? """; int updated = jdbcTemplate.update(sql, lockKey, lockOwner, LocalDateTime.now().plus(lockTimeout), lockOwner); return updated > 0; } @Override public boolean tryLock(long time, TimeUnit unit) throws InterruptedException { long endTime = System.currentTimeMillis() + unit.toMillis(time); while (System.currentTimeMillis() < endTime) { if (tryLock()) { return true; } Thread.sleep(100); } return false; } @Override public void unlock() { String sql = """ DELETE FROM kinesis_lock WHERE lock_key = ? AND owner = ? """; jdbcTemplate.update(sql, lockKey, lockOwner); } @Override public Condition newCondition() { throw new UnsupportedOperationException("暂不支持Condition"); } } }
步骤4:配置Spring Cloud Stream
修改application.yml,指定消费组、自定义的Checkpoint Store和Lock Registry:
spring: cloud: aws: credentials: sts: web-identity-token-file: <你的token文件路径> role-arn: <你的角色ARN> role-session-name: RoleSessionName region: static: <你的AWS区域> dualstack-enabled: false stream: bindings: input-in-0: destination: test-test.tst.v1 content-type: text/json group: kinesis-consumer-group # 必须指定消费组,Checkpoint与消费组绑定 kinesis: binder: auto-create-stream: false min-shard-count: 1 checkpoint-store: postgresCheckpointStore # 自定义Checkpoint Store的Bean名称 lock-registry: postgresLockRegistry # 自定义Lock Registry的Bean名称 checkpoint-mode: RECORD # 可选:RECORD(每条记录后Checkpoint) / BATCH(批量后Checkpoint)
步骤5:简化消费者代码
现在无需手动处理Checkpoint逻辑,Binder会自动调用你的CheckpointStore保存最新序列号,重启时自动从上次中断处恢复:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.Message; import java.util.function.Consumer; @Configuration public class KinesisConsumerBinder { @Bean public Consumer<Message<String>> input() { return message -> { System.out.println("Data from Kinesis: " + message.getPayload()); // 业务处理逻辑,处理成功后Binder自动完成Checkpoint }; } }
核心优势
- 完全复用现有PostgreSQL资源,无需额外部署DynamoDB
- 依托Spring Cloud Stream原生机制,避免手动处理Checkpoint和锁的重复逻辑
- 支持多实例分布式部署,保证分片消费的独占性
- 应用重启时自动从上次Checkpoint位置恢复消费
内容的提问来源于stack exchange,提问作者Jerrin Paul
相关产品推荐
相关产品推荐

