【问题标题】:Kinesis Java consumer - keepalive/heartbeat?Kinesis Java 消费者 - 保持活动/心跳?
【发布时间】:2017-06-21 05:20:34
【问题描述】:

我编写了一个使用 Kinesis 主题的 Java 服务。只要数据不太频繁,服务就可以很好地启动/运行并愉快地使用。如果有 > 60-90 分钟的间隙,它将停止消耗。不会发出任何错误,但随后的数据会在 Kinesis 中排队并等待服务重新启动。

Kinesis 是否需要在这些安静期间发送某种心跳或保活消息?

我查看了配置 (KinesisClientLibConfiguration) 并没有看到任何明显的东西。希望这不会需要每小时循环连接。


编辑:

KinesisClientLibConfiguration kinesisClientLibConfiguration = 
    new KinesisClientLibConfiguration(config.getString("appname"),
        config.getString("kinesis/stream_name"),
        kinesisCredentialsProvider, localProvider, 
        localProvider, workerId);

kinesisClientLibConfiguration.withInitialPositionInStream(
      InitialPositionInStream.TRIM_HORIZON);

编辑:

我设法找到了一些错误输出 - 其中很多:

com.amazonaws.services.kinesis.clientlibrary.lib.worker.ProcessTask 调用 严重:ShardId shardId-000000000000:捕获异常: com.amazonaws.services.kinesis.model.AmazonKinesisException:请求中包含的安全令牌已过期(服务:AmazonKinesis;状态代码:400;错误代码:ExpiredTokenException;请求 ID:cdb95cb6-23bb-0067-9c7b-1ad1125d7b7e)

这些消息在应用启动后正好 60 分钟开始。发现 this 引用说“在到期前 5 分钟刷新”。鉴于我在此调用中有两种凭据(一种用于 kinesis,另一种用于 dynamodb/cloud watch),我将尝试使用计时器到 .refresh()。

【问题讨论】:

  • 您可能想要显示您的KinesisClientLibConfiguration 代码。我猜你的消费者配置是因为 Shard 达到了END_OF_SHARD?你也可以在 dynamodb 中检查你的 consumerTable 吗?

标签: java amazon-kinesis


【解决方案1】:

(这是有效的)

注意:这涉及两个凭据来源 - 一个本地和一个远程。本地的用于 DynamoDB 和 CloudWatch。遥控器适用于 Kinesis。

AWSCredentialsProvider localProvider = new ClasspathPropertiesFileCredentialsProvider("credentials");

STSAssumeRoleSessionCredentialsProvider stsRoleCredentials = new STSAssumeRoleSessionCredentialsProvider.Builder(
  config.getString("kinesis/arn"), config.getString("kinesis/role_session_name"))
  .withExternalId(config.getString("kinesis/external_id")).build();

KinesisClientLibConfiguration kinesisClientLibConfiguration = new KinesisClientLibConfiguration(
  config.getString("appname"),
  config.getString("kinesis/stream_name"),
  stsRoleCredentials, localProvider, localProvider, workerId);

如果您使用 CredentialsProvider,它将为您刷新令牌。经过一番探索,我发现我正在进行的.refresh() 调用是为了清空函数。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-06-20
    • 2019-08-21
    • 2018-08-20
    • 2015-10-23
    • 1970-01-01
    相关资源
    最近更新 更多