DataGen 커넥터

DataGen 커넥터

DataGen 커넥터는 Flink 파이프라인의 입력 데이터를 생성할 수 있는 Source 구현을 제공합니다. Kafka와 같은 외부 시스템에 접근할 수 없을 때 로컬 개발이나 데모에 유용합니다. DataGen 커넥터는 내장되어 있어 추가 의존성이 필요하지 않습니다.

출처: 문서

본문

사용법

DataGeneratorSource는 N개의 데이터 포인트를 병렬로 생성합니다. 소스는 시퀀스를 병렬 소스 서브태스크 수만큼의 병렬 하위 시퀀스로 나눕니다. 사용자가 제공한 GeneratorFunctionLong 타입의 "index" 값을 공급하여 데이터 생성 과정을 구동합니다.

그런 다음 GeneratorFunction은 (하위) 시퀀스의 Long 값을 임의 데이터 타입의 생성된 이벤트로 매핑하는 데 사용됩니다. 예를 들어 다음 코드는 ["Number: 0", "Number: 1", ... , "Number: 999"] 레코드 시퀀스를 생성합니다.

GeneratorFunction<Long, String> generatorFunction = index -> "Number: " + index;
long numberOfRecords = 1000;

DataGeneratorSource<String> source =
        new DataGeneratorSource<>(generatorFunction, numberOfRecords, Types.STRING);

DataStreamSource<String> stream =
        env.fromSource(source,
        WatermarkStrategy.noWatermarks(),
        "Generator Source");

요소의 순서는 병렬도에 따라 달라집니다. 각 하위 시퀀스는 순서대로 생성됩니다. 따라서 병렬도가 1로 제한되면 "Number: 0"부터 "Number: 999"까지 순서대로 하나의 시퀀스를 생성합니다.

속도 제한 (Rate Limiting)

DataGeneratorSource는 속도 제한을 기본적으로 지원합니다. 다음 코드는 전체 소스 속도(모든 소스 서브태스크 합산)가 초당 100 이벤트를 초과하지 않는 String 값의 스트림을 생성합니다.

GeneratorFunction<Long, String> generatorFunction = index -> "Number: " + index;

DataGeneratorSource<String> source =
        new DataGeneratorSource<>(
             generatorFunction,
             Long.MAX_VALUE,
             RateLimiterStrategy.perSecond(100),
             Types.STRING);

checkpoint당 레코드 수 제한 등의 추가 속도 제한 전략은 RateLimiterStrategy에서 찾을 수 있습니다.

유계성 (Boundedness)

이 소스는 항상 유계(bounded)입니다. 다만 실용적으로는 레코드 수를 Long.MAX_VALUE로 설정하면 사실상 무계(unbounded) 소스가 됩니다(끝에 도달하지 않음). 유한 시퀀스의 경우 사용자는 BATCH 실행 모드로 애플리케이션을 실행하는 것을 고려할 수 있습니다.

참고 사항

참고: DataGeneratorSourceGeneratorFunction의 출력이 입력에 대해 결정적일 때, 즉 같은 Long 숫자를 공급하면 항상 같은 출력을 생성한다는 조건에서 at-least-once 및 end-to-end exactly-once 처리 보장을 가진 Flink 작업을 구현하는 데 사용할 수 있습니다.

참고: 생성된 이벤트와 사용자 정의 WatermarkStrategy를 기반으로 소스에서 바로 결정적 워터마크를 생성하는 것도 가능합니다.

더 알아보기 (Learn more)