【发布时间】:2014-11-11 12:16:27
【问题描述】:
好的,我将从一个详细的用例开始并解释我的问题:
- 我使用第 3 方网络分析平台,该平台利用 AWS Kinesis 流将数据从客户端传递到最终目的地 - Kinesis 流;
- 网络分析平台使用 2 个流:
- 数据收集器流(单个分片流);
- 第二个流,用于丰富来自收集器流(单个分片流)的原始数据;最重要的是,此流使用
TRIM_HORIZON迭代器类型使用来自第一个流的原始数据;
- 我使用 AWS Java SDK 使用流中的数据,具体使用
GetShardIteratorRequest类; - 我目前正在开发提取类,所以这是同步完成的,这意味着我只在编译我的类时才使用数据;
- 该类的效果出人意料,尽管有些事情我无法理解,特别是关于如何从流中使用数据以及每种迭代器类型的含义;
我的问题是我检索到的数据不一致,里面没有时间逻辑。
-
当我使用
提供分片中的第一个序列号时AT_SEQUENCE_NUMBER并使用.getSequenceNumberRange().getStartingSequenceNumber();
...作为``,我没有得到所有记录。同样,
AFTER_SEQUENCE_NUMBER; - 当我使用
LATEST时,我得到的结果为零; - 当我使用
TRIM_HORIZON时,它应该是有意义的,但它似乎不能正常工作。它曾经为我提供数据,然后我添加了新的“事件”(记录到最终流中)并且我收到了零记录。神秘。
我的问题是:
- 如何安全地使用流中的数据,而不必担心丢失记录?
- 是否有替代
ShardIteratorRequest的方法? - 如果有,我如何才能“浏览”流并查看其中的内容以供调试参考?
-
TRIM_HORIZON方法缺少什么?
提前致谢,我真的很想了解更多有关 Kinesis 流中数据消耗的信息。
【问题讨论】:
-
我也有类似的问题 - 尽管对我来说,我在每次迭代中都会得到重复的记录(同时使用 AT_SEQUENCE_NUMBER 和 FROM_SEQUENCE_NUMBER),尽管每次响应都使用 NextShardIterator 值。文档在这个问题上有些神秘......我也很想知道“未修剪”是什么意思(w.r.t TRIM_HORIZON)。
-
为了记录,我在此期间做了一些不同的事情——我使用了一个现有的 Scala 消费者,它不断地监听流,然后为了我的目的将它移植回纯 Java。这是 Scala 应用程序,最初由 SnowPlow github.com/snowplow/kinesis-example-scala-consumer 开发
-
可悲的是,我对 java 不友好.....!我只是希望有一个与语言无关的明确指南,说明如何确保幂等性和 100% 的记录“覆盖”,同时允许消费者重启、崩溃等。如果我们必须保存并检查 SequenceNumber,这似乎否定了 Kinesis 的目的所有先前提取的记录,以确保没有重复。我确定我错过了一些东西.......
-
您尝试过亚马逊自己的库吗? github.com/awslabs/amazon-kinesis-connectors github.com/awslabs/amazon-kinesis-client 这些库(尤其是连接器)处理所有繁琐的事情,例如查明检查点、继续处理分片等。
-
我在使用没有 KCL 的 JSON api 时遇到了类似的问题。我想获取最后一条记录作为检查点。 LATEST 给了我一个空数组。 TRIM_HORIZON 目前给了我 8 条记录。我可以遍历所有记录(可能是数千个)以获得最后一个,但这似乎很荒谬。 latest 应该如何工作?无论 KCL 在做什么,它都应该使用相同的 API,说“使用 KCL”并不能回答问题,它的检查点应该只基于这个 API 和存储的结果。
标签: java scala amazon-web-services amazon-kinesis