【问题标题】:Can I "branch" stream into many and write them in parallel in pyspark?我可以将流“分支”成许多流并在 pyspark 中并行写入它们吗?
【发布时间】:2021-07-14 01:52:07
【问题描述】:

我在 pyspark 中接收 Kafka 流。目前我正在按一组字段对其进行分组并将更新写入数据库:

df = spark \
        .readStream \
        .format("kafka") \
        .option("kafka.bootstrap.servers", config["kafka"]["bootstrap.servers"]) \
        .option("subscribe", topic)

...

df = df \
        .groupBy("myfield1") \
        .agg(
            expr("count(*) as cnt"),
            min(struct(col("mycol.myfield").alias("mmm"), col("*"))).alias("minData")
        ) \
        .select("cnt", "minData.*") \
        .select(
            col("...").alias("..."),
            ...
            col("userId").alias("user_id")

query = df \
        .writeStream \
        .outputMode("update") \
        .foreachBatch(lambda df, epoch: write_data_frame(table_name, df, epoch)) \
        .start()

query.awaitTermination()

我可以在中间使用相同的链并创建另一个分组

df2 = df \
        .groupBy("myfield2") \
        .agg(
            expr("count(*) as cnt"),
            min(struct(col("mycol.myfield").alias("mmm"), col("*"))).alias("minData")
        ) \
        .select("cnt", "minData.*") \
        .select(
            col("...").alias("..."),
            ...
            col("userId").alias("user_id")

并将其输出并行写入不同的位置?

在哪里打电话给writeStream和awaitTermination?

【问题讨论】:

    标签: pyspark apache-kafka spark-structured-streaming


    【解决方案1】:

    是的,您可以将 Kafka 输入流分支到任意数量的流式查询中。

    您需要考虑以下几点:

    1. query.awaitTermination 是一种阻塞方法,这意味着您在此方法之后编写的任何代码都不会在此 query 终止之前执行。
    2. 每个“分支”流式查询都将并行运行,在每个 writeStream 调用中定义检查点位置很重要。

    总体而言,您的代码需要具有以下结构:

    df = spark \
            .readStream \
            .format("kafka") \
            .option("kafka.bootstrap.servers", config["kafka"]["bootstrap.servers"]) \
            .option("subscribe", topic) \
            .[...]
    
    # note that I changed the variable name to "df1"
    df1 = df \
        .groupBy("myfield1") \
        .[...]
    
    df2 = df \
        .groupBy("myfield2") \
        .[...]
    
    
    query1 = df1 \
            .writeStream \
            .outputMode("update") \
            .option("checkpointLocation", "/tmp/checkpointLoc1") \
            .foreachBatch(lambda df, epoch: write_data_frame(table_name, df1, epoch)) \
            .start()
    
    query2 = df2 \
            .writeStream \
            .outputMode("update") \
            .option("checkpointLocation", "/tmp/checkpointLoc2") \
            .foreachBatch(lambda df, epoch: write_data_frame(table_name, df2, epoch)) \
            .start()
    
    spark.streams.awaitAnyTermination
    

    补充一句:在您显示的代码中,您将覆盖df,因此df2 的派生可能无法获得预期的结果。

    【讨论】:

    • 我觉得更好的方法是将记录序列化为字节数组或字符串,合并数据帧,然后写入
    • 不确定我是否理解您的观点,但我理解 OP 提出的问题,即如何让两个不同的聚合数据流来自同一来源。我的印象是生成的两个并行流将具有不同的架构,因此将它们合并起来会很困难。
    猜你喜欢
    • 1970-01-01
    • 2020-08-19
    • 2018-06-04
    • 2016-09-21
    • 2018-08-09
    • 2021-08-19
    • 2012-02-10
    • 2015-11-22
    • 1970-01-01
    相关资源
    最近更新 更多