【问题标题】:take sample form streaming dataframe采取样本形式流数据帧
【发布时间】:2020-04-02 13:06:50
【问题描述】:

我正在尝试将一个函数(适用于常规 spark 数据帧)应用于流数据。在我应用这个函数之前,我需要对给定的数据使用 .rdd.takeSample() 但当然这不适用于流数据帧。

我使用以下结构化流代码获取流数据:

dsraw = spark \
            .readStream \
            .format("kafka") \
            .option("kafka.bootstrap.servers", "192.168.99.100:9092") \
            .option("subscribe", "topic") \
            .option("startingOffsets", "earliest") \
            .load()

ds = dsraw.selectExpr("CAST(value AS STRING)")

我的数据是一组随机数,形式为 {'number': 1} 等。理想情况下,我想将从该流中读取的所有数字放入一个数据帧中并返回。

有没有办法将流数据帧转换为 spark 数据帧或 rdd?如果没有,takeSample 是否有替代方法?

【问题讨论】:

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


    【解决方案1】:

    一种方法是将流数据写入内存,
    然后使用spark sql:

    创建数据框/rdd
    dsraw = spark \
                .readStream \
                .format("kafka") \
                .option("kafka.bootstrap.servers", "192.168.99.100:9092") \
                .option("subscribe", "topic") \
                .option("startingOffsets", "earliest") \
                .load()
    
    ds = dsraw.selectExpr("CAST(value AS STRING)")
    
    kafka_value_df = ds.selectExpr("CAST(value AS STRING)")
    output_query = kafka_value_df.writeStream \
                          .queryName("numbers") \
                          .format("memory") \
                          .start()
    output_query.awaitTermination(10)
    
    value_df = spark.sql("select * from numbers")  # df
    
    value_rdd = value_df.rdd  # rdd
    

    我不确切知道您的原始数据是什么格式(仅{'number': 1} 是不够的信息),您可能必须使用mapjson.loads 取决于您的数据以获得所需的df 格式/rdd.

    【讨论】:

      猜你喜欢
      • 2018-02-01
      • 1970-01-01
      • 1970-01-01
      • 2017-09-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-26
      • 2020-07-04
      相关资源
      最近更新 更多