【问题标题】:Transform a specific group of records from MongoDB从 MongoDB 转换一组特定的记录
【发布时间】:2018-07-30 19:16:30
【问题描述】:

我有一个定期触发的批处理作业,它将数据写入 MongoDB。这项工作大约需要 10 分钟,之后我想接收这些数据并使用 Apache Flink 进行一些转换(映射、过滤、清理......)。记录之间存在一些依赖关系,这意味着我必须一起处理它们。例如,我喜欢转换客户 ID 为 45666 的最新批处理作业中的所有记录。结果将是一条聚合记录。

是否有任何最佳实践或方法可以做到这一点,而无需自己实现所有内容(从最新工作中获取不同的客户 ID,为每个客户选择记录和转换,标记转换后的客户等......)?

我无法流式传输它,因为我必须将多条记录一起转换,而不是一一转换。

目前我正在使用 Spring Batch、MongoDB、Kafka 并考虑使用 Apache Flink。

【问题讨论】:

  • 只想指出,即使您必须同时转换多条记录,使用有状态流式处理也可能有意义。 Flink 将让您保持记录状态,直到您拥有产生结果所需的所有部分。当然,这可能是一个好主意,也可能不是一个好主意,这取决于您的其他要求。
  • 我只确定在读取整个源文件之前没有关于组的更多信息,并且可能在 10 到 35 GB 之间。

标签: apache-flink


【解决方案1】:

可以想象,您可以将 MongoDB 更改流连接到 Flink,并将其用作您描述的任务的基础。涉及 10-35 GB 数据这一事实并不排除使用 Flink 流,因为您可以将 Flink 配置为在其状态无法放入堆时溢出到磁盘。

不过,在得出结论认为这是一种明智的做法之前,我希望更好地了解情况。

【讨论】:

  • 假设我有一个 30 GB 的 CSV 平面文件。在该文件的第一行,有一个来自 Customer_1 的条目,其中 [type=B, cus_id=1]。在文件的最后一行,有来自 Customer_1 的第二个条目 [type=A, zip=1234, cus_id=1]。意味着该客户的第二个条目比第一个条目具有更多和不同的信息。我不想在不知道会有第二个/不同的条目的情况下处理和插入第一个条目。我将处理/插入 [type=A, zip=1234, cus_id=1] -> 像 upsert 一样。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-11-09
  • 1970-01-01
  • 2014-03-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多