【问题标题】:Persisting Kinesis messages to S3 in Parquet format以 Parquet 格式将 Kinesis 消息持久保存到 S3
【发布时间】:2019-05-28 13:55:42
【问题描述】:

我有 Kinesis 流,我的应用每秒向其写入大约 10K 条原始格式的消息。

我想以 parquet 格式将这些消息保存到 S3。为了事后方便搜索,我需要按用户 ID 字段对我的数据进行分区,这是消息的一部分。

目前,我有一个由 Kinesis 事件触发的 lambda 函数。它接收多达 10K 条消息,按用户 ID 对它们进行分组,然后将这些文件以 parquet 格式写入 S3。

我的问题是这个 lambda 函数生成的文件非常小,大约 200KB,而我想创建大约 200MB 的文件以获得更好的查询性能(我使用 AWS Athena 查询这些文件)。

天真的方法是编写另一个 lambda 函数来读取这些文件并将它们合并(汇总)到一个大文件中,但我觉得我错过了一些东西,必须有更好的方法来做到这一点。

我想知道是否应该按照this 问题中的描述使用 Spark。

【问题讨论】:

  • 在查看建议的编辑时请小心。您最近批准了this,它的改进为零,并且根据评论,仅用于测试编辑。
  • @jhpratt,好的,但是您的评论与此线程有什么关系?

标签: apache-spark amazon-s3 parquet amazon-kinesis


【解决方案1】:

也许您可以使用 AWS 提供的两项额外服务:

AWS Kinesis Data Analytics 使用来自 Kinesis Stream 的数据并对您的数据(组、过滤器等)生成 SQL 分析。在此处查看更多信息:https://aws.amazon.com/kinesis/data-analytics/

AWS Kinesis Firehose 在 Kinesis Data Analytics 之后插入。使用此服务,我们可以每 X 分钟或每 Y MB 在 s3 上创建一个 parquet 文件,其中包含到达的数据。在此处查看更多信息:https://docs.aws.amazon.com/firehose/latest/dev/what-is-this-service.html

第二种方法是使用 Spark Structured Streaming。因此,您可以从 AWS Kinesis Stream 读取数据,过滤不可用的数据并导出到 s3,如下所述: https://databricks.com/blog/2017/08/09/apache-sparks-structured-streaming-with-amazon-kinesis-on-databricks.html

P.S.:这个例子展示了如何输出到本地文件系统,但你可以将其更改为 s3 位置。

【讨论】:

  • 谢谢。至于 Kinesis Firehouse,我认为我不能使用它对数据进行分区。据我了解,它会读取所有消息并将它们写入一个文件,而我需要按用户 ID 对它们进行分组。
  • Kinesis firehose 可以使用来自 Kinesis 流的数据,并每隔 X 分钟或每 Y MB 输出 parquet 文件(它可以生成许多文件)。按 ID 分组对您来说真的很重要吗?如果答案是肯定的,则必须使用 Spark Structured Stream。
  • 使用 AWS Kinesis 数据分析通过云解决方案更新答案
猜你喜欢
  • 2020-10-06
  • 2014-12-07
  • 1970-01-01
  • 2021-08-15
  • 1970-01-01
  • 1970-01-01
  • 2021-08-09
  • 2016-02-06
相关资源
最近更新 更多