PySpark如何将单行列表转换为DataFrame并写入Parquet表?
解决方案
问题根源
你的代码错误在于将Spark Column对象(如F.lit()、F.current_timestamp())直接放入列表,尝试用parallelize或.toDF()转换为DataFrame。parallelize仅支持本地Python数据(如整数、字符串、datetime),无法序列化Spark Column对象;而.toDF()需要传入行/元组列表,不是Column表达式列表。此外,你错误地将F.lit(datetime.datetime.now())赋值给start_time,导致后续重复包裹Column,引发类型错误。
正确实现方式
方法一:直接用spark.select()创建DataFrame
这是最简洁的方式,直接通过Column表达式构建单行DataFrame:
from pyspark.sql import functions as F from pyspark.sql.session import SparkSession import datetime # 初始化Spark会话(必须先执行) spark = SparkSession.builder.appName("SingleRowParquet").getOrCreate() tablename = 'table' start_time = datetime.datetime.now() # 本地起始时间(Python datetime类型) count_1 = 0 count_2 = 0 count_3 = 0 # 用select直接构建DataFrame,每个Column指定别名 df = spark.select( F.lit(start_time).alias('start_time'), F.lit(tablename).alias('table'), F.lit(count_1).alias('count_1'), F.lit(count_2).alias('count_2'), F.lit(count_3).alias('count_3'), F.current_timestamp().alias('end_time') ) # 写入Parquet表(mode可选append/overwrite等) df.write.parquet("/path/to/your/parquet/directory", mode="append")
方法二:先创建静态DataFrame,再添加动态列
如果需要先处理静态数据,再添加Spark生成的动态列(如current_timestamp()):
from pyspark.sql import functions as F from pyspark.sql.session import SparkSession from pyspark.sql.types import StructType, StructField, TimestampType, StringType, IntegerType import datetime spark = SparkSession.builder.appName("SingleRowParquet").getOrCreate() tablename = 'table' start_time = datetime.datetime.now() count_1 = 0 count_2 = 0 count_3 = 0 # 显式定义Schema(可选,但能确保类型正确) schema = StructType([ StructField("start_time", TimestampType(), nullable=False), StructField("table", StringType(), nullable=False), StructField("count_1", IntegerType(), nullable=False), StructField("count_2", IntegerType(), nullable=False), StructField("count_3", IntegerType(), nullable=False) ]) # 用本地数据创建初始DataFrame data = [(start_time, tablename, count_1, count_2, count_3)] df = spark.createDataFrame(data, schema=schema) # 添加end_time列(由Spark生成当前时间戳) df = df.withColumn("end_time", F.current_timestamp()) # 写入Parquet表 df.write.parquet("/path/to/your/parquet/directory", mode="append")
关键注意事项
- 必须先初始化Spark会话(
SparkSession.builder),否则无法创建DataFrame。 F.lit()仅用于包裹本地字面量(如整数、字符串、datetime),不能包裹Spark Column对象。F.current_timestamp()是Spark内置函数,会在DataFrame执行时生成集群节点的当前时间戳。
内容的提问来源于stack exchange,提问作者wmbff
相关产品推荐
相关产品推荐

