- 38
- 0
批次间隔为10s, 窗口大小为20s, 步长为10s, 这样每个window应该有2个批次的数据,但是我用DStream.foreachRDD()每次只执行一次,按我理解因为有2个批次数据应该执行两次,但实际测试下来无论window中有多少batch都是只调用一次
如何辨别出每一个批次的数据呢?比如第一个批次执行某个操作,第二个批次执行另一种操作,但他们都在同一个窗口中
- 共 0 条
- 全部回答
-
笑看往事如花 普通会员 1楼
在Spark Stream中,你可以使用
filter和map函数来依次遍历同一个窗口中每个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'。
- 扫一扫访问手机版
回答动态

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器更新之后。服务器里面有部分玩家要重新创建角色是怎么回事啊?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题函数计算不同地域的是不能用内网吧?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题ARMS可以创建多个应用嘛?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题在ARMS如何申请加入公测呀?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题前端小程序接入这个arms具体是如何接入监控的,这个init方法在哪里进行添加?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器刚到期,是不是就不能再导出存档了呢?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器的游戏版本不兼容 尝试更新怎么解决?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器服务器升级以后 就链接不上了,怎么办?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器转移以后服务器进不去了,怎么解决?预计能赚取 0积分收益

- 神奇的四哥:发布了悬赏问题阿里云幻兽帕鲁服务器修改参数后游戏进入不了,是什么情况?预计能赚取 0积分收益
- 回到顶部
- 回到顶部
