【发布时间】: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","result":"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