代码之家  ›  专栏  ›  技术社区  ›  sansari

添加检查点程序时不使用记录

  •  1
  • sansari  · 技术社区  · 8 年前

    我有以下配置KinesisMessageDrivenChannelAdapter,当我删除 dynamoDbMetaDataStore 作为检查点,消息被正确地接收,但是当我添加它时,记录总是空的。 我调试了代码 KinesisMessageDrivenChannelAdapter.processTask() 第776行(版本2.0.0.m2)返回空记录。

    更新:

    public DynamoDbMetaDataStore dynamoDbMetaDataStore() {
        String url = consumerClientProperties.getDynamoDB().getUrl();
        final AmazonDynamoDBAsync amazonDynamoDB = AmazonDynamoDBAsyncClientBuilder.standard()
            .withEndpointConfiguration(new EndpointConfiguration(
                url,
                Regions.fromName(awsRegion).getName()))
            .withClientConfiguration(new ClientConfiguration()
                .withMaxErrorRetry(consumerClientProperties.getDynamoDB().getRetries())
                .withConnectionTimeout(consumerClientProperties.getDynamoDB().getConnectionTimeout())).build();
        DynamoDbMetaDataStore dynamoDbMetaDataStore = new DynamoDbMetaDataStore(amazonDynamoDB, "consumer-test");
        return dynamoDbMetaDataStore;
      }
    
      public KinesisMessageDrivenChannelAdapter kinesisInboundChannel(
          AmazonKinesis amazonKinesis, String[] streamNames) {
        KinesisMessageDrivenChannelAdapter adapter =
            new KinesisMessageDrivenChannelAdapter(amazonKinesis, streamNames);
        adapter.setConverter(null);
        adapter.setOutputChannel(kinesisReceiveChannel());
        adapter.setCheckpointStore(dynamoDbMetaDataStore());
        adapter.setConsumerGroup(consumerClientProperties.getName());
        adapter.setCheckpointMode(CheckpointMode.manual);
        adapter.setListenerMode(ListenerMode.record);
        adapter.setStartTimeout(10000);
        adapter.setDescribeStreamRetries(1);
        adapter.setConcurrency(10);
        return adapter;
      }
    

    谢谢你

    1 回复  |  直到 8 年前
        1
  •  0
  •   Artem Bilan    8 年前

    我建议你用最新的 2.0.0.BUILD-SNAPSHOT .

    已经有一个选项类似于:

    /**
     * Specify a {@link LockRegistry} for an exclusive access to provided streams.
     * This is not used when shards-based configuration is provided.
     * @param lockRegistry the {@link LockRegistry} to use.
     * @since 2.0
     */
    public void setLockRegistry(LockRegistry lockRegistry) {
    

    你需要注射 DynamoDbLockRegistry 为了更好的检查点管理。

    为此,您还需要添加此依赖项:

    compile("com.amazonaws:dynamodb-lock-client:1.0.0")
    

    在这个过程中过滤确实可能会有一些问题 M2 然而…

    推荐文章