代码之家  ›  专栏  ›  技术社区  ›  JDev

Spark结构化流使用多个查询的用例

  •  0
  • JDev  · 技术社区  · 5 年前

    我需要从多个Kafka主题(基于Avro)中进行流式传输,并将它们放入Greenplum中,只需对有效载荷进行小幅修改。

    Kaka主题被定义为配置文件中的列表,每个Kafka主题都有一个目标表。

    我正在寻找一个Spark Structured应用程序和配置文件中的更新,以收听新主题或停止。听这个话题。

    我正在寻求帮助,因为我对使用单个查询和多个查询感到困惑:

    val query1 = df.writeStream.start()
    val query2 = df.writeStream.start()
    
    spark.streams.awaitAnyTermination()
    

    df.writeStream.start().awaitAnyTermination()
    

    在何种用例下,应在单个查询上使用多个查询

    0 回复  |  直到 5 年前
        1
  •  1
  •   Shane    5 年前

    显然,您可以使用正则表达式模式来消费来自不同kafka主题的数据。

    假设你有主题名称,如“top-ingesion1”、“top-iangesion2”,那么你可以创建一个正则表达式模式,用于消费所有以“*摄入”结尾的主题的数据。

    一旦以正则表达式模式的格式创建了新主题,spark将自动开始从新创建的主题流式传输数据。

    参考: [https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html#consumer-缓存]

    您可以使用此参数指定缓存超时。 “火花。卡夫卡。消费者。缓存。超时”。

    来自spark文档:

    spark.kafka.consumer.cache.timeout-最短时间 消费者可能会在游泳池里闲着,然后才有资格被驱逐 被驱逐者。

    假设你有多个接收器,在那里你从kafka读取数据,并将其写入两个不同的位置,如hdfs和hbase,那么你可以将应用程序逻辑分支到两个writeStreams中。

    如果接收器(Greenplum)支持批处理操作模式,那么您可以查看spark结构化流中的forEachBatch()函数。它将允许我们为这两个操作重用相同的batchDF。

    参考: [https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html#consumer-缓存]

    推荐文章