【问题标题】:Pick up changes in json files that are being read by pyspark readstream?获取 pyspark readstream 正在读取的 json 文件中的更改?
【发布时间】:2023-02-22 04:04:49
【问题描述】:

我有 json 文件,其中每个文件描述一个特定的实体,包括它的状态。我试图通过使用 readStream 和 writeStream 将它们拉入 Delta。这对新文件非常有效。这些 json 文件经常更新(即状态更改、添加 cmets、添加历史项等)。更改后的 json 文件不会被 readStream 拉入。我认为这是因为 readStream 不重新处理项目。有没有解决的办法?

我正在考虑的一件事是更改我对 json 的初始写入以向文件名添加时间戳,以便它成为与流不同的记录(无论如何我已经必须在我的 writeStream 中进行重复数据删除),但我是尝试不修改正在编写 json 的代码,因为它已经在生产中使用。

理想情况下,我想找到类似 Cosmos Db 的 changeFeed 功能的东西,但用于读取 json 文件。

有什么建议么?

谢谢!

【问题讨论】:

    标签: pyspark spark-structured-streaming azure-synapse delta-lake


    【解决方案1】:

    Spark Structured Streaming 不支持此功能 - 文件处理后不会再次处理。

    最接近您的要求的只存在于Databricks' Autoloader - 它有选项cloudFiles.allowOverwrites option 允许重新处理修改后的文件。

    附言如果您对文件源 (https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#input-sources) 使用 cleanSource 选项,那么它可能会重新处理文件,但我不是 100% 确定。

    【讨论】:

    • 感谢 Alex 提供的信息,它确实可以帮助我实现其他目的,但不能用于此。 cleanSource 将帮助我清除文件并避免它产生必须在大目录上跟踪目录列表的开销,但即使如此检查点也会阻止重新处理。对于 Autoloader,这听起来不错,但我没有提到我在 Azure Synapse 中工作,而 Autoloader 是 Databricks 的一项功能。我还有其他 Databricks 项目,我可以在其中进行研究。
    • 是的,我注意到您正在使用 Synapse,但我不知道非 Databricks 实现的此类功能
    【解决方案2】:

    运气好吗?我正在寻找类似的功能

    【讨论】:

    • Mohit Israni,请不要添加我也是作为答案。它实际上并没有提供问题的答案。如果您有不同但相关的问题,请ask它(如果它有助于提供上下文,请参考此问题)。如果你对这个具体问题感兴趣,你可以upvote它,留下comment,或者一旦你有足够的reputation就开始bounty
    猜你喜欢
    • 2023-04-07
    • 1970-01-01
    • 2020-01-08
    • 2021-06-30
    • 1970-01-01
    • 2021-12-17
    • 2021-09-28
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多