Avro 포맷

Avro 포맷 (Avro format)

Flink는 Apache Avro를 내장으로 지원합니다. 즉, Avro 스키마를 기반으로 Avro 데이터를 쉽게 읽고 쓸 수 있습니다. Flink의 직렬화 프레임워크는 Avro 스키마로부터 생성된 클래스를 처리할 수 있습니다. Avro 포맷을 사용하려면 빌드 자동화 도구(Maven, SBT 등)를 쓰는 프로젝트에서 다음 의존성이 필요합니다.

출처: 문서

본문

<dependency>
  <groupId>org.apache.flink</groupId>
  <artifactId>flink-avro</artifactId>
  <version>2.3.0</version>
</dependency>

PyFlink 작업에서 Avro 포맷을 사용하려면 PyFlink JAR이 필요합니다. PyFlink에서 JAR을 사용하는 방법은 Python 의존성 관리 문서를 참고하세요.

Avro 파일에서 데이터를 읽으려면 AvroInputFormat을 지정해야 합니다.

예시:

AvroInputFormat<User> users = new AvroInputFormat<User>(in, User.class);
DataStream<User> usersDS = env.createInput(users);

여기서 User는 Avro가 생성한 POJO입니다. Flink는 이러한 POJO의 문자열 기반 키 선택도 지원합니다. 예:

usersDS.keyBy("name");

GenericData.Record 타입을 Flink에서 사용하는 것은 가능하지만 권장하지 않습니다. 레코드가 전체 스키마를 포함하고 있어 데이터가 매우 크고 따라서 느릴 수 있기 때문입니다.

Flink의 POJO 필드 선택은 Avro로 생성된 POJO에서도 동작합니다. 단, 필드 타입이 생성된 클래스에 올바르게 기록된 경우에만 가능합니다. 필드가 Object 타입이면 조인이나 그룹 키로 사용할 수 없습니다. Avro에서 필드를 {"name": "type_double_test", "type": "double"},처럼 지정하면 잘 동작하지만, 단일 필드 UNION 타입({"name": "type_double_test", "type": ["double"]},)으로 지정하면 Object 타입 필드가 생성됩니다. nullable 타입({"name": "type_double_test", "type": ["null", "double"]},)을 지정하는 것은 가능합니다.

Python 작업의 경우 Avro 파일에서 읽으려면 Avro 스키마를 정의해야 하며, 요소는 일반 Python 객체가 됩니다. 예:

schema = AvroSchema.parse_string("""
{
    "type": "record",
    "name": "User",
    "fields": [
        {"name": "name", "type": "string"},
        {"name": "favoriteNumber",  "type": ["int", "null"]},
        {"name": "favoriteColor", "type": ["string", "null"]}
    ]
}
""")

env = StreamExecutionEnvironment.get_execution_environment()
ds = env.create_input(AvroInputFormat(AVRO_FILE_PATH, schema))

def json_dumps(record):
    import json
    return json.dumps(record)

ds.map(json_dumps).print()

더 알아보기 (Learn more)