두 형식은 지향이 다르다. Parquet 은 컬럼 지향으로 분석 조회에 맞고, Avro 는 행 지향으로 스트리밍 전송과 스키마 진화에 맞는다. Kafka 로 흘려보내야 해서 Avro 가 필요한 경우가 대표적이다.
pyarrow 로 읽고 fastavro 로 쓴다. Arrow 에는 Avro 쓰기 기능이 없으므로 스키마를 직접 만들어 넘겨야 한다.
pip install pyarrow fastavro
Arrow 타입과 Avro 타입을 대응시킨다. 자동으로 해 주는 함수는 없다.
import pyarrow as pa
ARROW_TO_AVRO = {
pa.bool_(): 'boolean',
pa.int32(): 'int',
pa.int64(): 'long',
pa.float32(): 'float',
pa.float64(): 'double',
pa.string(): 'string',
pa.binary(): 'bytes',
pa.date32(): {'type': 'int', 'logicalType': 'date'},
pa.timestamp('ms'): {'type': 'long', 'logicalType': 'timestamp-millis'},
pa.timestamp('us'): {'type': 'long', 'logicalType': 'timestamp-micros'},
}
def avro_type(field: pa.Field):
for arrow_type, mapped in ARROW_TO_AVRO.items():
if field.type.equals(arrow_type):
break
else:
raise TypeError(f'매핑하지 않은 타입: {field.name} {field.type}')
# NULL 을 허용하는 컬럼은 union 으로 감싸고 기본값을 null 로 둔다
if field.nullable:
return ['null', mapped]
return mapped
def build_schema(arrow_schema: pa.Schema, name='Record', namespace='local'):
fields = []
for f in arrow_schema:
entry = {'name': f.name, 'type': avro_type(f)}
if f.nullable:
entry['default'] = None
fields.append(entry)
return {'type': 'record', 'name': name, 'namespace': namespace, 'fields': fields}
decimal 은 Avro 에서 bytes 에 logicalType: decimal 과 precision · scale 을 붙여야 한다. 위 표에 넣지 않았으니 필요하면 따로 다룬다.
파일 전체를 메모리에 올리면 큰 파일에서 터진다. Parquet 은 row group 단위로 나뉘어 있으므로 그 단위로 읽어 배치로 쓴다.
import pyarrow.parquet as pq
import fastavro
def parquet_to_avro(src: str, dst: str, codec: str = 'snappy'):
pf = pq.ParquetFile(src)
schema = fastavro.parse_schema(build_schema(pf.schema_arrow))
with open(dst, 'wb') as out:
first = True
for batch in pf.iter_batches(batch_size=10_000):
records = batch.to_pylist()
if first:
fastavro.writer(out, schema, records, codec=codec)
first = False
else:
# 두 번째 배치부터는 기존 파일에 이어 붙인다
fastavro.writer(out, schema, records, codec=codec)
if __name__ == '__main__':
parquet_to_avro('input.parquet', 'output.avro')
fastavro.writer 는 파일 객체와 레코드 목록을 받는다. 같은 열린 파일 객체에 반복 호출하면 이어 붙는다. 이어쓰기를 따로 명시하려면 파일을 'a+b' 로 열고 fastavro.writer(out, None, records) 처럼 스키마를 None 으로 주면 기존 파일의 스키마를 읽어 쓴다.
codec 은 null(무압축) · deflate · snappy · zstandard 를 쓸 수 있다. snappy 는 python-snappy, zstandard 는 zstandard 패키지가 따로 필요하다.
import fastavro
with open('output.avro', 'rb') as f:
reader = fastavro.reader(f)
print(reader.writer_schema)
for i, rec in enumerate(reader):
print(rec)
if i >= 4:
break
수십 GB 를 이 방식으로 돌리는 것은 적절하지 않다. Spark 로 처리한다. 스키마 변환도 알아서 한다.
spark.read.parquet('hdfs:///raw/orders') \
.write.format('avro').save('hdfs:///stage/orders_avro')
Avro 데이터소스는 별도 패키지로 제공되므로 --packages org.apache.spark:spark-avro_2.12:<SPARK_VERSION> 을 붙여 제출한다.