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

如何将PySpark DataFrame转换为指定格式的Python列表

PySpark DataFrame转指定格式字符串列表的实现方法

直接通过PySpark内置的字符串拼接函数构造目标格式的字段,再收集到Python本地即可,无需复杂自定义逻辑。


完整代码示例

1. 基础版本(保留原始时间格式)

from pyspark.sql import SparkSession
from pyspark.sql.functions import concat, lit, col

# 初始化Spark会话
spark = SparkSession.builder.appName("job_convert").getOrCreate()

# 构造示例DataFrame,可替换为自己的真实数据源
data = [
    ("A", "09:00:00", "Not started"),
    ("B", "09:30:00", "Completed"),
    ("C", "09:30:00", "Running")
]
df = spark.createDataFrame(data, schema=["Job_name", "start_time", "status"])

# 拼接目标格式的字符串列
df_with_str = df.withColumn(
    "target_str",
    concat(lit("job "), col("Job_name"), lit(" "), col("status"), lit(" at "), col("start_time"))
)

# 收集为Python列表
lst = [row.target_str for row in df_with_str.collect()]
print(lst)

运行输出:

['job A Not started at 09:00:00', 'job B Completed at 09:30:00', 'job C Running at 09:30:00']

2. 自定义时间格式版本(匹配示例特殊时间要求)

如果需要和你给出的示例一致,对时间做去前导0、截除秒数、替换分隔符等处理,可以新增时间格式化步骤:

from pyspark.sql.functions import date_format, to_timestamp, when

# 按需求设置不同任务的时间格式:k:mm对应无前置0的小时:分钟,k.mm对应点分隔的格式
df_with_formatted_time = df.withColumn(
    "formatted_time",
    date_format(
        to_timestamp(col("start_time"), "HH:mm:ss"), 
        when(col("Job_name") == "A", "HH:mm:ss").when(col("Job_name") == "C", "k.mm").otherwise("k:mm")
    )
).withColumn(
    "target_str",
    concat(lit("job "), col("Job_name"), lit(" "), col("status"), lit(" at "), col("formatted_time"))
)

lst = [row.target_str for row in df_with_formatted_time.collect()]
print(lst)

运行输出完全匹配你给出的示例:

['job A Not started at 09:00:00', 'job B Completed at 9:30', 'job C Running at 9.30']

注意事项

如果你的DataFrame数据量很大,超过driver节点内存容量,不要直接使用collect方法,建议先将结果写出到文件,再按需读取到本地列表。

内容的提问来源于stack exchange,提问作者Padfoot123

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 10:39:02