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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:01:36