【问题标题】:Spark Streaming with schema-less data具有无模式数据的 Spark Streaming
【发布时间】:2017-12-20 08:15:11
【问题描述】:

我们目前有一个数据管道设置,我们正在使用 Logstash 从单个 Kafka 主题读取原始数据并将其写入 ElasticSearch。
本主题中的数据是JSON格式的,但是每一行可以属于完全不同的业务领域,所以可能有完全不同的schema。例如:

记录 1:"{"id":1,"model":"model2","updated":"2017-01-1T00:00:00.000Z","domain":"A"}

记录 2:"{"id":"some_compound_key","re​​sult":"PASS","domain":"B"}

您可以看到,不仅架构不同,而且实际上是冲突的(id 在第一条记录中是整数,而在第二条记录中是字符串)。

只有两个保证——每条记录都是有效的 JSON 记录,每条记录都有一个“域”字段。即使具有相同域值的记录有时也可能具有不同的架构。

我们现在需要在这些数据通过管道时对其进行丰富和转换(而不是稍后使用 ETL 进行),并且我们正在研究几种方法来完成它。需要注意的是,由于数据没有统一的模式,因此需要逐行进行转换:

1) 继续使用 Logstash - 可以使用一组 Logstash 过滤器和条件来为每个域建模我们需要的转换管道。
由于 Logstash 在运行时会定期重新加载配置,因此维护和部署也很容易,因此要更改/添加转换逻辑,我们只需在 conf 目录中放置一个新的配置文件。
然而,缺点是很难使用 Logstash 从外部来源丰富数据。

2) 使用 Kafka 流 - 这似乎是一个显而易见的选择,因为它与 Kafka 很好地集成,允许连接来自多个流(或外部源)的数据并且没有架构要求 - 很容易逐行转换数据。
这里的缺点是很难在运行时修改转换逻辑 - 我们需要重新编译和重新部署应用程序,或者用一些 API 包装它,以便在运行时生成和编译 Java 代码,或者其他一些复杂的解决方案。

3) 使用 Spark 流式处理 - 我们已经在使用 Spark 进行批处理,因此如果我们也可以将其用于流式处理以使我们的堆栈尽可能简单,那就太好了。
但是,我不确定 Spark 是否甚至可以支持没有单一架构的流数据,也不确定是否可以按行执行转换。
我见过的所有示例(以及我们自己使用 Spark 批处理的经验)都假设数据具有明确定义的架构,这不是我们的用例。

任何人都可以阐明我们需要的 Spark Streaming(或结构化流)是否可以实现,还是我们应该坚持使用 Logstash / Kafka Streams?

【问题讨论】:

    标签: apache-spark spark-streaming apache-kafka-streams


    【解决方案1】:

    免责声明:我是 Kafka Streams 的积极贡献者。

    我对 Logstash 不熟悉,但从您的描述来看,它似乎是最没有吸引力的解决方案。

    关于 Spark Streaming。即使我不是它的忠实粉丝,我相信你可以用它做你想做的处理。根据我的理解,结构化流式处理不起作用,因为它需要一个固定的模式,但 Spark 流式处理应该更灵活。然而,与 Kafka Streams 相比,使用 Spark Streaming 并不会让它变得更简单(但很可能更难)。我没有在生产中运行 Spark Streaming 的个人经验,但我听到了很多关于不稳定等的抱怨。

    关于您指出的 Kafka Streams 的“缺点”。 (1) 我不确定您为什么需要代码生成等以及 (2),为什么这在 Spark Streaming 中会有所不同?你需要在这两种情况下编写你的转换逻辑,如果你想改变它,你需要重新部署。我也相信,通过“滚动反弹”更新 Kafka Streams 应用程序更容易,并且与您需要在其间停止处理的 Spark Streaming 相比允许零停机时间。

    了解您想要执行的“运行时代码修改”会有所帮助,以便在此处给出更详细的答案。

    【讨论】:

    • Kafka Streams 和 Spark Streaming 在可管理性方面没有区别 - 两者都需要重新编译应用程序并重新部署它 - 这是我们希望尽可能避免并找到支持的方法例如,在运行时添加新的 lambda 函数。我同意 Kafka Streams 更灵活——因为每个流应用程序都是独立的,我们可以通过添加更多容器化实例轻松地向上/向下扩展它们,而使用 Spark 则无法做到这一点。您在滚动反弹方面也提出了一个很好的观点 - 我没有考虑过:-) 谢谢!
    • “Spark 结构化流需要架构”是有争议的。我们在 prod 上部署了与模式完全无关的管道。
    • @TroubleShooter 你是如何实现的?
    • @VictorKironde,如果您可以消除反序列化负载的需要(请注意,您可以使用 JSON 库检查负载中的任何字段)任何火花管道都可以与模式无关。对于 JSON,我们使用 io.gatling.jsonpath 来查找消息中的特定字段,从 Spark 的角度来看,这已经足够了。
    猜你喜欢
    • 2018-04-04
    • 2021-11-24
    • 1970-01-01
    • 2020-03-19
    • 2021-04-26
    • 2018-02-02
    • 2015-08-07
    • 2016-12-13
    • 2017-05-20
    相关资源
    最近更新 更多