如何将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
相关产品推荐
相关产品推荐

