【问题标题】:Read from Kafka and write to hdfs in parquet从 Kafka 读取并写入 parquet 中的 hdfs
【发布时间】:2018-01-31 07:58:14
【问题描述】:

我是 BigData 生态系统的新手,并且刚刚起步。

我已经阅读了几篇关于使用 spark 流式传输阅读 kafka 主题的文章,但想知道是否可以使用 spark 作业而不是流式传输从 kafka 中读取? 如果是的话,你们能否帮我指出一些可以让我开始的文章或代码sn-ps。

我的第二部分问题是以 parquet 格式写入 hdfs。 一旦我从 Kafka 阅读,我假设我会有一个 rdd。 将此 rdd 转换为数据帧,然后将数据帧写入 parquet 文件。 这是正确的方法吗?

任何帮助表示赞赏。

谢谢

【问题讨论】:

    标签: hadoop apache-spark apache-kafka hdfs parquet


    【解决方案1】:

    要从 Kafka 读取数据并将其以 Parquet 格式写入 HDFS,使用 Spark Ba​​tch 作业而不是流式传输,您可以使用 Spark Structured Streaming

    Structured Streaming 是基于 Spark SQL 引擎构建的可扩展和容错流处理引擎。您可以像表达对静态数据的批处理计算一样表达您的流计算。 Spark SQL 引擎将负责以增量和连续的方式运行它,并随着流数据的不断到达而更新最终结果。您可以使用 Scala、Java、Python 或 R 中的 Dataset/DataFrame API 来表示流聚合、事件时间窗口、流到批处理连接等。计算在同一个优化的 Spark SQL 引擎上执行。最后,系统通过检查点和预写日志确保端到端的精确一次容错保证。简而言之,结构化流式处理提供了快速、可扩展、容错、端到端的一次性流处理,而用户无需对流式处理进行推理。

    它带有 Kafka 作为内置 Source,即我们可以从 Kafka 轮询数据。它与 Kafka 代理版本 0.10.0 或更高版本兼容。

    为了以批处理模式从 Kafka 中提取数据,您可以为定义的偏移范围创建 Dataset/DataFrame。

    // Subscribe to 1 topic defaults to the earliest and latest offsets
    val df = spark
      .read
      .format("kafka")
      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
      .option("subscribe", "topic1")
      .load()
    df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .as[(String, String)]
    
    // Subscribe to multiple topics, specifying explicit Kafka offsets
    val df = spark
      .read
      .format("kafka")
      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
      .option("subscribe", "topic1,topic2")
      .option("startingOffsets", """{"topic1":{"0":23,"1":-2},"topic2":{"0":-2}}""")
      .option("endingOffsets", """{"topic1":{"0":50,"1":-1},"topic2":{"0":-1}}""")
      .load()
    df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .as[(String, String)]
    
    // Subscribe to a pattern, at the earliest and latest offsets
    val df = spark
      .read
      .format("kafka")
      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
      .option("subscribePattern", "topic.*")
      .option("startingOffsets", "earliest")
      .option("endingOffsets", "latest")
      .load()
    df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .as[(String, String)]
    

    源中的每一行都有以下架构:

    | Column           | Type          |
    |:-----------------|--------------:|
    | key              |        binary |
    | value            |        binary |
    | topic            |        string |
    | partition        |           int |
    | offset           |          long |
    | timestamp        |          long |
    | timestampType    |           int |
    

    现在,要将数据以 parquet 格式写入 HDFS,可以编写以下代码:

    df.write.parquet("hdfs://data.parquet")
    

    有关 Spark Structured Streaming + Kafka 的更多信息,请参阅以下指南 - Kafka Integration Guide

    希望对你有帮助!

    【讨论】:

    • 这个答案有用吗?
    • 感谢 Himanshu,这很有帮助。似乎这需要 Spark 2.2,在 2.0 等较低版本的 spark 中有没有其他方法可以做到这一点。
    【解决方案2】:

    关于这个话题,你已经有了几个很好的答案。

    只是想强调一下 - 小心直接流入镶木地板。 当 parquet 行组大小足够大(为简单起见,您可以说文件大小应为 64-256Mb 的顺序),利用字典压缩、布隆过滤器等(一个 parquet 文件可以有多个行块,并且通常每个文件中确实有多个行块;尽管行块不能跨越多个拼花文件)

    如果您直接流式传输到 parquet 表,那么您很可能会得到一堆很小的 parquet 文件(取决于 Spark Streaming 的小批量大小和数据量)。查询此类文件可能非常慢。例如,Parquet 可能需要读取所有文件的标题以协调模式,这是一个很大的开销。如果是这种情况,您将需要一个单独的进程,例如,作为一种解决方法,读取旧文件并将它们“合并”写入(这不是简单的文件级合并,进程将实际上需要读入所有 parquet 数据并溢出更大的 parquet 文件)。

    这种解决方法可能会破坏数据“流式传输”的最初目的。你也可以在这里查看其他技术——比如 Apache Kudu、Apache Kafka、Apache Druid、Kinesis 等,它们可以在这里更好地工作。

    更新:自从我发布了这个答案后,这里现在有了一个新的强者 - Delta Lakehttps://delta.io/如果你习惯parquet,你会发现Delta很有吸引力(实际上,Delta是建立在parquet层+元数据之上的)。 Delta Lake 提供:

    Spark 上的 ACID 事务:

    • 可序列化的隔离级别确保读取器永远不会看到不一致的数据。
    • 可扩展的元数据处理:利用 Spark 的分布式处理能力轻松处理数十亿文件的 PB 级表的所有元数据。
    • 流和批处理统一:Delta Lake 中的表是批处理表以及流源和接收器。流式数据摄取、批量历史回填、交互式查询都可以开箱即用。
    • 架构强制执行:自动处理架构变化,以防止在提取期间插入不良记录。
    • 时间旅行:数据版本控制支持回滚、完整的历史审计跟踪和可重复的机器学习实验。
    • Upserts 和删除:支持合并、更新和删除操作,以支持复杂的用例,例如更改数据捕获、缓慢变化的维度 (SCD) 操作、流式更新插入等。

    【讨论】:

      【解决方案3】:

      使用 Kafka 流。 SparkStreaming 用词不当(它是底层的小批量,至少高达 2.2)。

      https://eng.verizondigitalmedia.com/2017/04/28/Kafka-to-Hdfs-ParquetSerializer/

      【讨论】:

      猜你喜欢
      • 2018-10-24
      • 2021-10-19
      • 2021-07-22
      • 2018-12-18
      • 2016-06-29
      • 2018-01-18
      • 2019-04-12
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多