如何在PySpark DataFrame中合并同ID的send与receive行至含对应列的单行
Solution to Combine PySpark Rows by ID with Pivoted Action Columns
Got it, let's work through this problem. You need to merge paired 'receive' and 'send' rows for each ID into a single row, with dedicated columns for each message type, plus the earliest timestamp for each ID. Here's a straightforward way to do it in PySpark:
Step-by-Step Approach
- Group by ID: Since each ID has exactly one 'receive' and one 'send' entry, grouping by
IDlets us process each pair together. - Capture the earliest timestamp: Use
min(MSG_DT)to get the first (earliest) timestamp for each ID group. - Map messages to new columns: Use conditional logic to pull the
MESSAGEvalue whereACTION_CDis 'receive' or 'send', and assign them to newRECEIVEandSENDcolumns.
PySpark Code Implementation
Assuming your original DataFrame is named df, here are two reliable methods:
Method 1: Direct Aggregation with Conditional Logic
This is the most concise approach for your specific case (only two action types):
from pyspark.sql import functions as F # Group by ID and compute required fields result_df = df.groupBy("ID") \ .agg( # Get the earliest timestamp for the ID F.min("MSG_DT").alias("MSG_DT"), # Extract MESSAGE where ACTION_CD is 'receive' F.first(F.when(F.col("ACTION_CD") == "receive", F.col("MESSAGE"))).alias("RECEIVE"), # Extract MESSAGE where ACTION_CD is 'send' F.first(F.when(F.col("ACTION_CD") == "send", F.col("MESSAGE"))).alias("SEND") ) # Show the result result_df.show(truncate=False)
Method 2: Using Pivot (Scalable for More Action Types)
If you might add more ACTION_CD values later, using pivot is more flexible. We first get the earliest timestamp, then join with a pivoted version of the message data:
from pyspark.sql import functions as F # First, get the earliest MSG_DT per ID timestamps_df = df.groupBy("ID") \ .agg(F.min("MSG_DT").alias("MSG_DT")) # Pivot the ACTION_CD to get message columns per action type pivoted_messages_df = df.groupBy("ID") \ .pivot("ACTION_CD") \ .agg(F.first("MESSAGE")) # Join the two DataFrames to combine timestamps and messages result_df = timestamps_df.join(pivoted_messages_df, on="ID") # Show the result result_df.show(truncate=False)
Expected Output
Both methods will produce the exact DataFrame you're looking for:
+------------------------------------+------------------------+-------+-------------------+ |ID |MSG_DT |RECEIVE|SEND | +------------------------------------+------------------------+-------+-------------------+ |d2636151-b95e-4845-8014-0a113c381ff9|2022-08-07T21:24:54.552Z|Oi |Olá! | |4241224b-9ba5-4eda-8e16-7e3aeaacf164|2022-08-07T21:25:06.565Z|4 |Certo | |bd46c6fb-1315-4418-9943-2e7d3151f788|2022-08-07T21:25:30.408Z|1 |⭐️*Antes de você ir| |14da8519-6e4c-4edc-88ea-e33c14533dd9|2022-08-07T21:25:52.798Z|788884 |Agora | +------------------------------------+------------------------+-------+-------------------+
内容的提问来源于stack exchange,提问作者Sergio Pantano
相关产品推荐
相关产品推荐

