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