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

如何将Parquet文件转换为Avro文件?求解Pandas转换报错问题

解决Parquet转Avro并获取Avro Schema的问题

首先,你遇到的AttributeError: 'DataFrame' object has no attribute 'write'错误是因为df.write.format("avro")是Spark DataFrame的API,而你用Pandas读取得到的是Pandas DataFrame,两者API不兼容。

下面提供几种可行的转换方案,附带获取Avro Schema的方法:

方案1:使用PyArrow直接转换(推荐,和你当前用的引擎一致)

PyArrow支持直接读写Parquet和Avro文件,无需转成Pandas DataFrame,效率更高:

import pyarrow.parquet as pq
import pyarrow.avro as avro

# 读取Parquet文件为PyArrow Table
parquet_table = pq.read_table('local_output-00000-of-00001.parquet')

# 将Table写入Avro文件
avro.write_avro('deserialize.avro', parquet_table)

# 读取Avro文件并提取Schema
with open('deserialize.avro', 'rb') as f:
    avro_reader = avro.AvroFileReader(f)
    avro_schema = avro_reader.schema
    print(avro_schema)

方案2:Pandas结合fastavro实现转换

如果需要基于Pandas DataFrame操作,可以使用fastavro库:

  1. 先安装依赖:
pip install fastavro
  1. 转换代码:
import pandas as pd
from fastavro import writer, reader, parse_schema

# 读取Parquet文件到Pandas DataFrame
df = pd.read_parquet('local_output-00000-of-00001.parquet', engine='pyarrow')

# 将DataFrame转为字典列表格式
records = df.to_dict('records')

# 从DataFrame字段类型推断Avro Schema(可根据需求自定义)
def infer_avro_schema(df):
    schema_fields = []
    for col_name, dtype in df.dtypes.items():
        # 映射Pandas类型到Avro类型
        if pd.api.types.is_integer_dtype(dtype):
            avro_type = 'long'
        elif pd.api.types.is_float_dtype(dtype):
            avro_type = 'double'
        elif pd.api.types.is_string_dtype(dtype):
            avro_type = 'string'
        elif pd.api.types.is_bool_dtype(dtype):
            avro_type = 'boolean'
        # 扩展支持更多类型(如日期、数组等)
        else:
            avro_type = 'string'  # 默认 fallback 到string
        schema_fields.append({'name': col_name, 'type': avro_type})
    return {'type': 'record', 'name': 'ParquetToAvroRecord', 'fields': schema_fields}

# 解析并验证Schema
avro_schema = parse_schema(infer_avro_schema(df))

# 写入Avro文件
with open('deserialize.avro', 'wb') as out_file:
    writer(out_file, avro_schema, records)

# 读取Avro文件并获取Schema
with open('deserialize.avro', 'rb') as in_file:
    avro_reader = reader(in_file)
    print(avro_reader.schema)

方案3:使用Spark(适合大数据量场景)

如果你原本想使用Spark的API,需要确保环境中安装了Spark Avro包,正确代码如下:

from pyspark.sql import SparkSession

# 初始化SparkSession(需提前安装Spark及Avro包)
spark = SparkSession.builder \
    .appName("ParquetToAvro") \
    .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.5.0") \
    .getOrCreate()

# 读取Parquet文件为Spark DataFrame
df = spark.read.parquet('local_output-00000-of-00001.parquet')

# 写入Avro文件
df.write.format("avro").save("deserialize.avro")

# 读取Avro文件并获取Schema
avro_df = spark.read.format("avro").load("deserialize.avro")
print(avro_df.schema)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:15:18