You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

  1. Group by ID: Since each ID has exactly one 'receive' and one 'send' entry, grouping by ID lets us process each pair together.
  2. Capture the earliest timestamp: Use min(MSG_DT) to get the first (earliest) timestamp for each ID group.
  3. Map messages to new columns: Use conditional logic to pull the MESSAGE value where ACTION_CD is 'receive' or 'send', and assign them to new RECEIVE and SEND columns.

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.27 17:14:09