如何将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库:
- 先安装依赖:
pip install fastavro
- 转换代码:
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
相关产品推荐
相关产品推荐

