账号密码登录
微信安全登录
微信扫描二维码登录

登录后绑定QQ、微信即可实现信息互通

手机验证码登录
找回密码返回
邮箱找回 手机找回
注册账号返回
其他登录方式
分享
  • 收藏
    X
    spark-stream中如何依次遍历同一个window中每个batch的数据呢?
    38
    0

    批次间隔为10s, 窗口大小为20s, 步长为10s, 这样每个window应该有2个批次的数据,但是我用DStream.foreachRDD()每次只执行一次,按我理解因为有2个批次数据应该执行两次,但实际测试下来无论window中有多少batch都是只调用一次

    如何辨别出每一个批次的数据呢?比如第一个批次执行某个操作,第二个批次执行另一种操作,但他们都在同一个窗口中

    4
    打赏
    收藏
    点击回答
        全部回答
    • 0
    • 笑看往事如花 普通会员 1楼

      在Spark Stream中,你可以使用filtermap函数来依次遍历同一个窗口中每个batch的数据。以下是一个简单的示例:

      ```python from pyspark.sql import SparkSession

      创建SparkSession

      spark = SparkSession.builder.getOrCreate()

      假设你有一个名为'data'的DataFrame,其中包含'batch_id'和'data'两列

      df = spark.read.format('csv').option('header', 'true').load('path/to/data.csv')

      创建一个窗口变量

      window_id = df.columns[0]

      使用filter函数来过滤出在同一个窗口中的数据

      filtered_data = df.filter(lambda row: row[window_id] == 'batch1')

      使用map函数将过滤后的数据转换为新的列

      new_column = filtered_data.map(lambda row: row[1])

      打印结果

      new_column.show() ```

      在这个示例中,我们首先创建了一个SparkSession,然后从CSV文件中读取数据。然后,我们创建了一个窗口变量,并使用filter函数过滤出在同一个窗口中的数据。最后,我们使用map函数将过滤后的数据转换为新的列。

      注意,你需要将'path/to/data.csv'替换为你的实际数据文件路径。此外,你需要将'batch1'替换为你的实际数据中的'batch_id'。

    更多回答
    扫一扫访问手机版
    • 回到顶部
    • 回到顶部