【问题标题】:Reading Message from Kafka Topic and Dump it into HDFS从 Kafka 主题读取消息并将其转储到 HDFS
【发布时间】:2020-08-19 13:20:12
【问题描述】:

我正在尝试使用来自 Kafka 主题的数据,将其加载到数据集中,然后在加载到 Hdfs 之前执行过滤。

我能够从 kafka 主题中消费,将其加载到数据集中并保存为 HDFS 中的镶木地板文件,但无法执行过滤条件。你能分享一下在保存到hdfs之前执行过滤的方法吗? 我正在使用 Java 和 Spark 从 kafka 主题中消费。 我的部分代码是这样的:

DataframeDeserializer dataframe = new DataframeDeserializer(dataset);

 ds = dataframe.fromConfluentAvro("value", <your schema path>, <yourmap>, RETAIN_SELECTED_COLUMN_ONLY$.MODULE$);

StreamingQuery query = ds.coalesce(10)
                .writeStream()
                .format("parquet")
                .option("path", path.toString())
                .option("checkpointLocation", "<your path>")
                .trigger(Trigger.Once())
                .start();

【问题讨论】:

  • @Srinivas 的回答很好。我想我们可以学习kafka-connect-hdfs 的工作原理或使用这个库。

标签: java scala apache-spark apache-kafka hdfs


【解决方案1】:

我强烈推荐Kafka Connect,而不是重新发明轮子。 您只需要 HDFS Sink 连接器,它将数据从 Kafka 主题复制到 HDFS。

【讨论】:

    【解决方案2】:

    coalesce之前写过滤逻辑,即ds.filter().coalesce()

    
    DataframeDeserializer dataframe = new DataframeDeserializer(dataset);
    
     ds = dataframe.fromConfluentAvro("value", <your schema path>, <yourmap>, RETAIN_SELECTED_COLUMN_ONLY$.MODULE$);
    
    StreamingQuery query = 
                    ds
                    .filter(...) // Write your filter condition here
                    .coalesce(10)
                    .writeStream()
                    .format("parquet")
                    .option("path", path.toString())
                    .option("checkpointLocation", "<your path>")
                    .trigger(Trigger.Once())
                    .start();
    
    
    

    【讨论】:

    • 假设我有这样的过滤条件:Country='USA' 和 State='Texas',然后像这样的查询?: StreamingQuery query = ds.filter("County=='USA' AND St​​ate='Texas'") .coalesce(10) .writeStream() .format("parquet") .option("path", path.toString()) .option("checkpointLocation", "" ) .trigger(Trigger.Once()) .start();
    • 当然,Sri,我会根据你的建议修改代码,分享结果:)
    • 过滤条件按预期工作。感谢您的帮助。
    猜你喜欢
    • 2018-10-24
    • 1970-01-01
    • 2021-03-01
    • 2018-05-07
    • 1970-01-01
    • 2017-01-08
    • 2018-07-27
    • 2018-01-12
    • 1970-01-01
    相关资源
    最近更新 更多