PySpark中按diag_sid和vehicle_id分组打乱指定列行记录
在PySpark中分组打乱指定列对应行的方案
要实现按diag_sid和vehicle_id分组,打乱组内function_txt、item_txt、value_txt对应的行记录(保持三列对应关系,仅随机重排组内行顺序),可以用以下两种PySpark原生方案:
方案1:窗口函数(纯PySpark,推荐大数据场景)
利用窗口函数给每个分组内的行生成随机排序依据,再按该依据重排行:
from pyspark.sql import functions as F from pyspark.sql import Window # 定义分组窗口,按目标字段分组后随机排序 window_spec = Window.partitionBy("diag_sid", "vehicle_id").orderBy(F.rand()) # 添加随机行号,按行号重排后删除临时列 shuffled_df = ( df .withColumn("random_row_num", F.row_number().over(window_spec)) .orderBy("random_row_num") .drop("random_row_num") )
方案2:Pandas UDF(适合熟悉Pandas的场景)
通过分组后调用Pandas的排序逻辑实现:
from pyspark.sql import functions as F # 分组后用Pandas随机排序每组数据 shuffled_df = ( df .withColumn("rand_key", F.rand()) .groupBy("diag_sid", "vehicle_id") .applyInPandas( lambda group: group.sort_values("rand_key"), schema=df.schema.add("rand_key", "double") ) .drop("rand_key") )
问题说明
你之前使用的df.reindex(np.random.permutation(df.index))是Pandas的全局重排方法,只会打乱整个DataFrame的行顺序,无法实现分组内的行打乱,所以达不到预期效果。
输入示例
| diag_sid | vehicle_id | source | function_txt | item_txt | value_txt | date |
|---|---|---|---|---|---|---|
| 1711364453938960795 | W1KAH0FB1PF092903 | A | Bordnetzdaten_EZS222 | 22105E | 1A 00 18 CA 05 7A... | 2024-03-25 |
| 1711364453938960795 | W1KAH0FB1PF092903 | A | Bordnetzdaten_EZS223 | 221058 | 01 E4 FD E8 FD E0... | 2024-03-25 |
预期输出(组内行随机重排)
| diag_sid | vehicle_id | source | function_txt | item_txt | value_txt | date |
|---|---|---|---|---|---|---|
| 1711364453938960795 | W1KAH0FB1PF092903 | A | Bordnetzdaten_EZS223 | 221058 | 01 E4 FD E8 FD E0... | 2024-03-25 |
| 1711364453938960795 | W1KAH0FB1PF092903 | A | Bordnetzdaten_EZS222 | 22105E | 1A 00 18 CA 05 7A... | 2024-03-25 |
内容的提问来源于stack exchange,提问作者mb_eng
相关产品推荐
相关产品推荐

