You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 15:45:24