【问题标题】:Firehose datapipeline limitationsFirehose 数据管道限制
【发布时间】:2019-08-23 22:15:10
【问题描述】:

我的用例如下: 我有 JSON 数据进来,需要以镶木地板格式存储在 S3 中。到目前为止一切顺利,我可以在 Glue 中创建一个模式并将“DataFormatConversionConfiguration”附加到我的 firehose 流中。但是数据来自不同的“主题”。每个主题都有一个特定的“模式”。据我了解,我将不得不创建多个 firehose 流,因为一个流只能有一个模式。但是我有成千上万个这样的主题,其中有非常大量的高吞吐量数据传入。创建这么多的 firehose 资源看起来不太可行 (https://docs.aws.amazon.com/firehose/latest/dev/limits.html)

我应该如何构建我的管道。

【问题讨论】:

    标签: amazon-web-services bigdata amazon-kinesis-firehose data-pipeline


    【解决方案1】:

    您可以:

    • 要求升级您的 Firehose 限制,并使用 1 个 Firehose/流 + 添加 Lambda 转换以将数据转换为通用架构 - IMO 成本效益不高,但您应该看到您的负载。

    • 为每个 Kinesis 数据流创建一个 Lambda,将每个事件转换为由单个 Firehose 管理的架构,最后可以使用 Firehose API https://docs.aws.amazon.com/firehose/latest/APIReference/API_PutRecord.html 将事件直接发送到您的 Firehose 流(请参阅“Q:如何将数据添加到我的 Amazon Kinesis Data Firehose 传输流中?”此处为 https://aws.amazon.com/kinesis/data-firehose/faqs/) - 而且,请检查之前的成本,因为即使您的 Lambda 是“按需”调用的,您可能会在很长一段时间内调用其中的很多一段时间。

    • 使用其中一种数据处理框架(Apache Spark、Apache Flink 等)并以 1 小时为一批次从 Kinesis 读取数据,每次从您上次终止时开始 --> 使用可用的接收器转换数据并以 Parquet 格式写入。框架使用检查点的概念并将最后处理的偏移量存储在外部存储中。现在,如果您每小时重新启动它们,它们将开始直接从上次看到的条目中读取数据。 - 它可能具有成本效益,特别是如果您考虑使用 Spot 实例。另一方面,它需要比之前的 2 个解决方案更多的编码,并且显然可能有更高的延迟。

    希望对您有所帮助。您能否就选择的解决方案提供反馈?

    【讨论】:

    • 嘿,我们继续使用 Flink 和 Kinesis 解决方案。我们在 kinesis 流上运行它。设法使用单个 flink 接收器,通过向代码添加一些自定义扩展来动态检测新模式和主题:)
    • @Dexter 您是否将 Kinesis Data Streams 或 Kinesis Data Analytics 与 Flink 一起使用?我的问题和你的问题完全一样。还有错误是如何表示的,正如我在 Firehose 中发现的那样,如果有一些错误的验证,它没有给出正确的错误消息。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-03
    • 1970-01-01
    • 1970-01-01
    • 2021-10-10
    • 1970-01-01
    相关资源
    最近更新 更多