빌딩 블록

빌딩 블록 (Building Blocks)

DataStream API V2에서 DataStream, Partitioning, ProcessFunction은 가장 기본적인 요소예요. 각각 데이터 스트림의 유형, 데이터가 어떻게 파티셔닝되는지, 데이터 스트림에서 어떤 연산/처리를 수행하는지를 나타내요.

출처: 문서

본문

참고: DataStream API V2는 기존 DataStream API를 점진적으로 대체하기 위한 새로운 API 집합이에요. 현재 실험 단계이며 프로덕션에는 완전히 사용할 수 없어요.

DataStream, Partitioning, ProcessFunction은 DataStream API V2의 가장 기본적인 요소이며, 각각 다음을 나타내요.

  • 데이터 스트림의 유형이 무엇인지
  • 데이터가 어떻게 파티셔닝되는지
  • 데이터 스트림에서 연산/처리를 어떻게 수행하는지

이들은 또한 DataStream API가 제공하는 기본 프리미티브(primitive)의 핵심 부분이에요.

DataStream

스트림 위의 데이터 흐름은 여러 파티션으로 나뉠 수 있어요. 스트림에서 데이터가 어떻게 파티셔닝되는지에 따라 다음과 같은 범주로 나눠요.

  • Global Stream: 단일 파티션/병렬도(parallelism)를 강제하며, 데이터의 정확성은 이에 의존해요.
  • Keyed Partition Stream: 각 key가 파티션이고, 데이터가 속하는 파티션이 결정적이에요.
  • Non-Keyed Partition Stream: 각 parallelism이 파티션이고, 데이터가 속하는 파티션이 비결정적이에요.
  • Broadcast Stream: 각 파티션이 같은 데이터를 포함해요.

Partitioning

위에서 스트림과 그 파티셔닝 방식을 정의했어요. 다음 다룰 주제는 서로 다른 파티션 유형 사이를 변환하는 방법이에요. 이러한 변환을 partitioning이라고 불러요.

예를 들어 non-keyed partition stream은 KeyBy partitioning을 통해 keyed partition stream으로 변환될 수 있어요.

NonKeyedPartitionStream<Tuple<Integer, String>> stream = ...;
KeyedPartitionStream<Integer, String> keyedStream = stream.keyBy(record -> record.f0);

전체적으로 다음 네 가지 partitioning이 있어요.

  • KeyBy: 모든 데이터를 지정된 key에 따라 다시 파티셔닝해요.
  • Shuffle: 모든 데이터를 다시 파티셔닝하고 섞어요(shuffle).
  • Global: 모든 파티션을 하나로 병합해요.
  • Broadcast: 파티션이 다운스트림으로 데이터를 브로드캐스트하도록 강제해요.

구체적인 변환 관계는 다음 표에 나와 있어요. (교차된 칸은 지원되지 않거나 필요하지 않음을 나타내요.)

한 가지 주의할 점은: broadcast는 다른 입력과 결합해서만 사용할 수 있고, 다른 스트림으로 직접 변환될 수 없어요.

ProcessFunction

데이터 스트림을 가지면 여기에 연산을 적용할 수 있어요. DataStream에 대해 수행할 수 있는 연산들을 통틀어 Process Function이라고 불러요. 이는 데이터 스트림에서 모든 종류의 처리를 정의하는 유일한 진입점이에요.

ProcessFunction의 분류

입력/출력의 수에 따라 다음과 같이 분류돼요.

Partitioning 입력 수 출력 수
OneInputStreamProcessFunction 1 1
TwoInputNonBroadcastStreamProcessFunction 2 1
TwoInputBroadcastStreamProcessFunction 2 1
TwoOutputStreamProcessFunction 1 2

(더 많은 입력과 출력을 처리하려면 여러 process function을 결합해 달성할 수 있어요)

입력 중 하나가 broadcast 스트림인지 여부에 따라 두 가지 유형의 two-input process function이 있어요. DataStream에는 입력 스트림을 변환하거나, ProcessFunction을 통해 두 입력 스트림을 연결·변환하는 일련의 processconnectAndProcess 메서드가 있어요.

입력 및 출력 스트림 요구사항

다음 두 표는 OneInputStreamProcessFunctionTwoOutputStreamProcessFunction이 각각 지원하는 입력·출력 스트림 조합을 나열해요.

OneInputStreamProcessFunction의 경우:

입력 스트림 출력 스트림
Global Global
Keyed Keyed / NonKeyed
NonKeyed NonKeyed
Broadcast 지원 안 함

TwoOutputStreamProcessFunction의 경우:

입력 스트림 출력 스트림
Global Global + Global
Keyed Keyed + Keyed / Non-Keyed + Non-Keyed
NonKeyed NonKeyed + NonKeyed
Broadcast 지원 안 함

입력이 두 개인 경우는 조금 더 복잡해요. 다음 표는 어떤 스트림들이 서로 호환되는지, 그리고 어떤 유형의 스트림을 출력하는지를 나열해요. 교차(❎)는 지원되지 않음을 나타내요.

출력 Global Keyed NonKeyed Broadcast
Global Global
Keyed NonKeyed / Keyed NonKeyed / Keyed
NonKeyed NonKeyed NonKeyed
Broadcast NonKeyed / Keyed NonKeyed

프로세스 구성 (Config Process)

process function을 정의한 뒤에는 이 처리의 속성에 대해 몇 가지 구성을 하고 싶을 수 있어요. 예를 들어 process 연산의 parallelism과 이름을 설정하는 등이에요.

process/connectAndProcess의 반환값은 스트림이면서 동시에 이전 처리를 구성할 수 있게 해주는 핸들이에요. 구성은 withXXX라고 하는 여러 메서드로 해요. 예를 들어,

inputStream
  .process(func1) // do process 1
  .withName("my-process-func") // configure name for process 1
  .withParallelism(2) //  configure parallelism for process 1
  .process(func2) //  do further process 2

예시 (Example)

다음은 이러한 빌딩 블록을 사용해 flink job을 작성하는 방법을 보여주는 예시예요.

// create environment
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// create a stream from source
env.fromSource(someSource)
    // map every element x to x + 1.
    .process(new OneInputStreamProcessFunction<Integer, Integer>() {
                    @Override
                    public void processRecord(
                            Integer x,
                            Collector<Integer> output)
                            throws Exception {
                        output.collect(x + 1);
                    }
                })
    // If the sink does not support concurrent writes, we can force the stream to one partition
    .global()
    // sink the stream to some sink
    .toSink(someSink);
// execute the job
env.execute()

더 알아보기 (Learn more)