성능 튜닝

성능 튜닝 (Performance Tuning)

SQL 은 데이터 분석에 가장 널리 사용되는 언어입니다. Flink 의 Table API 와 SQL 은 사용자가 더 적은 시간과 노력으로 효율적인 스트림 분석 애플리케이션을 정의할 수 있게 해줍니다. 또한 Flink Table API 와 SQL 은 효과적으로 최적화되어 많은 쿼리 최적화와 튜닝된 operator 구현을 통합합니다. 그러나 모든 최적화가 기본적으로 활성화되지는 않으므로 일부 워크로드에서는 일부 옵션을 켜 성능을 개선할 수 있습니다.

출처: 문서

본문

SQL 은 데이터 분석에 가장 널리 사용되는 언어입니다. Flink 의 Table API 와 SQL 은 사용자가 더 적은 시간과 노력으로 효율적인 스트림 분석 애플리케이션을 정의할 수 있게 해줍니다. 또한 Flink Table API 와 SQL 은 효과적으로 최적화되어 많은 쿼리 최적화와 튜닝된 operator 구현을 통합합니다. 그러나 모든 최적화가 기본적으로 활성화되지는 않으므로 일부 워크로드에서는 일부 옵션을 켜 성능을 개선할 수 있습니다.

이 페이지에서 일부 경우에 큰 개선을 가져올 유용한 최적화 옵션과 스트리밍 집계, 정규 조인(regular join)의 내부 구조를 소개합니다.

이 페이지에서 언급하는 스트리밍 집계 최적화는 모두 현재 Group AggregationsWindow TVF Aggregations (Session Window TVF Aggregation 제외) 에서 지원됩니다.

MiniBatch 집계 (MiniBatch Aggregation)

기본적으로 그룹 집계 operator 는 입력 레코드를 하나씩 처리합니다. 즉, (1) 상태에서 accumulator 를 읽고, (2) accumulator 에 레코드를 누적/재시도하며, (3) accumulator 를 상태로 다시 쓰고, (4) 다음 레코드가 (1)부터 다시 처리합니다. 이 처리 패턴은 StateBackend 의 오버헤드를 증가시킬 수 있습니다(특히 RocksDB StateBackend 의 경우). 게다가 프로덕션에서 매우 흔한 데이터 왜곡(skew) 은 문제를 악화시켜 작업이 백프레셔 상황에 빠지기 쉽게 만듭니다.

mini-batch 집계의 핵심 아이디어는 집계 operator 내부의 버퍼에 입력 묶음(bundle)을 캐싱하는 것입니다. 입력 묶음이 처리되도록 트리거되면 키당 상태 접근은 한 번만 필요합니다. 이는 상태 오버헤드를 크게 줄이고 더 나은 처리량을 얻을 수 있습니다. 그러나 일부 레코드를 즉시 처리하지 않고 버퍼링하므로 지연 시간이 약간 늘어날 수 있습니다. 이는 처리량과 지연 시간 사이의 균형(trade-off)입니다.

MiniBatch 최적화는 그룹 집계에 대해 기본적으로 비활성화되어 있습니다. 활성화하려면 table.exec.mini-batch.enabled, table.exec.mini-batch.allow-latency, table.exec.mini-batch.size 옵션을 설정해야 합니다. 자세한 내용은 configuration 페이지를 참고하세요.

MiniBatch 최적화는 Window TVF Aggregation 을 위해 위의 구성과 무관하게 항상 활성화됩니다. Window TVF 집계는 레코드를 JVM Heap 대신 관리 메모리(managed memory) 에 버퍼링하므로 GC 오버로드나 OOM 문제의 위험이 없습니다.

다음 예제는 이러한 옵션을 활성화하는 방법을 보여줍니다.

Java:

// instantiate table environment
TableEnvironment tEnv = ...;

// access flink configuration
TableConfig configuration = tEnv.getConfig();
// set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true"); // enable mini-batch optimization
configuration.set("table.exec.mini-batch.allow-latency", "5 s"); // use 5 seconds to buffer input records
configuration.set("table.exec.mini-batch.size", "5000"); // the maximum number of records can be buffered by each aggregate operator task

Scala:

// instantiate table environment
val tEnv: TableEnvironment = ...

// access flink configuration
val configuration = tEnv.getConfig()
// set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true") // enable mini-batch optimization
configuration.set("table.exec.mini-batch.allow-latency", "5 s") // use 5 seconds to buffer input records
configuration.set("table.exec.mini-batch.size", "5000") // the maximum number of records can be buffered by each aggregate operator task

Python:

# instantiate table environment
t_env = ...

# access flink configuration
configuration = t_env.get_config()
# set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true") # enable mini-batch optimization
configuration.set("table.exec.mini-batch.allow-latency", "5 s") # use 5 seconds to buffer input records
configuration.set("table.exec.mini-batch.size", "5000") # the maximum number of records can be buffered by each aggregate operator task

Local-Global 집계

Local-Global 은 그룹 집계를 두 단계로 나누어 데이터 왜곡 문제를 해결하기 위해 제안되었습니다. 즉 먼저 업스트림에서 로컬 집계를 수행하고, 그 다음 다운스트림에서 전역 집계를 수행합니다. 이는 MapReduce 의 Combine + Reduce 패턴과 유사합니다. 예를 들어 다음 SQL 을 고려해 보세요:

SELECT color, sum(id)
FROM T
GROUP BY color

데이터 스트림의 레코드가 왜곡되어 있을 수 있으므로 집계 operator 의 일부 인스턴스가 다른 인스턴스보다 훨씬 많은 레코드를 처리해야 하며, 이는 핫스팟(hotspot)으로 이어집니다. 로컬 집계는 같은 키를 가진 일정량의 입력을 단일 accumulator 로 누적하는 데 도움이 될 수 있습니다. 전역 집계는 많은 수의 원시 입력 대신 축소된 accumulator 만 받습니다. 이는 네트워크 셔플과 상태 접근 비용을 크게 줄일 수 있습니다. 로컬 집계가 매번 누적하는 입력 수는 mini-batch 간격에 기반합니다. 즉 local-global 집계는 mini-batch 최적화가 활성화되어 있어야 합니다.

다음 예제는 local-global 집계를 활성화하는 방법을 보여줍니다.

Java:

// instantiate table environment
TableEnvironment tEnv = ...;

// access flink configuration
TableConfig configuration = tEnv.getConfig();
// set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true"); // local-global aggregation depends on mini-batch is enabled
configuration.set("table.exec.mini-batch.allow-latency", "5 s");
configuration.set("table.exec.mini-batch.size", "5000");
configuration.set("table.optimizer.agg-phase-strategy", "TWO_PHASE"); // enable two-phase, i.e. local-global aggregation

Scala:

// instantiate table environment
val tEnv: TableEnvironment = ...

// access flink configuration
val configuration = tEnv.getConfig()
// set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true") // local-global aggregation depends on mini-batch is enabled
configuration.set("table.exec.mini-batch.allow-latency", "5 s")
configuration.set("table.exec.mini-batch.size", "5000")
configuration.set("table.optimizer.agg-phase-strategy", "TWO_PHASE") // enable two-phase, i.e. local-global aggregation

Python:

# instantiate table environment
t_env = ...

# access flink configuration
configuration = t_env.get_config()
# set low-level key-value options
configuration.set("table.exec.mini-batch.enabled", "true") # local-global aggregation depends on mini-batch is enabled
configuration.set("table.exec.mini-batch.allow-latency", "5 s")
configuration.set("table.exec.mini-batch.size", "5000")
configuration.set("table.optimizer.agg-phase-strategy", "TWO_PHASE") # enable two-phase, i.e. local-global aggregation

Split Distinct 집계

Local-Global 최적화는 SUM, COUNT, MAX, MIN, AVG 같은 일반 집계의 데이터 왜곡을 제거하는 데 효과적입니다. 그러나 distinct 집계를 다룰 때는 그 성능이 만족스럽지 않습니다.

예를 들어 오늘 로그인한 고유 사용자 수를 분석하고 싶다고 가정해 보겠습니다. 다음 쿼리가 있을 수 있습니다:

SELECT day, COUNT(DISTINCT user_id)
FROM T
GROUP BY day

COUNT DISTINCT 는 distinct 키(예: user_id)의 값이 희소하면 레코드 줄이기에 좋지 않습니다. local-global 최적화가 활성화되어도 별 도움이 되지 않습니다. accumulator 가 여전히 거의 모든 원시 레코드를 포함하고 전역 집계가 병목이 되기 때문입니다(대부분의 무거운 accumulator 가 하나의 작업, 예: 같은 날짜에서 처리됨).

이 최적화의 아이디어는 distinct 집계(예: COUNT(DISTINCT col)) 를 두 수준으로 분할하는 것입니다. 첫 번째 집계는 그룹 키와 추가 버킷 키로 셔플됩니다. 버킷 키는 HASH_CODE(distinct_key) % BUCKET_NUM 으로 계산됩니다. BUCKET_NUM 은 기본적으로 1024이며 table.optimizer.distinct-agg.split.bucket-num 옵션으로 구성할 수 있습니다. 두 번째 집계는 원래 그룹 키로 셔플되고 SUM 을 사용해 서로 다른 버킷의 COUNT DISTINCT 값을 집계합니다. 같은 distinct 키는 같은 버킷에서만 계산되므로 변환은 동등합니다. 버킷 키는 그룹 키의 핫스팟 부담을 공유하는 추가 그룹 키 역할을 합니다. 버킷 키 덕분에 작업이 확장 가능해져 distinct 집계의 데이터 왜곡/핫스팟을 해결할 수 있습니다.

Split distinct 집계 후 위 쿼리는 자동으로 다음 쿼리로 다시 작성됩니다:

SELECT day, SUM(cnt)
FROM (
    SELECT day, COUNT(DISTINCT user_id) as cnt
    FROM T
    GROUP BY day, MOD(HASH_CODE(user_id), 1024)
)
GROUP BY day

참고: 위는 이 최적화의 혜택을 받을 수 있는 가장 간단한 예입니다. 그 외에도 Flink 는 더 복잡한 집계 쿼리의 분할을 지원합니다. 예를 들어 서로 다른 distinct 키를 가진 둘 이상의 distinct 집계(예: COUNT(DISTINCT a), SUM(DISTINCT b)), 다른 비-distinct 집계(예: SUM, MAX, MIN, COUNT) 와 함께 동작하는 경우 등입니다.

현재 split 최적화는 사용자 정의 AggregateFunction 을 포함하는 집계를 지원하지 않습니다.

다음 예제는 split distinct 집계 최적화를 활성화하는 방법을 보여줍니다.

Java:

// instantiate table environment
TableEnvironment tEnv = ...;

tEnv.getConfig()
  .set("table.optimizer.distinct-agg.split.enabled", "true");  // enable distinct agg split

Scala:

// instantiate table environment
val tEnv: TableEnvironment = ...

tEnv.getConfig
  .set("table.optimizer.distinct-agg.split.enabled", "true")  // enable distinct agg split

Python:

# instantiate table environment
t_env = ...

t_env.get_config().set("table.optimizer.distinct-agg.split.enabled", "true") # enable distinct agg split

Distinct 집계에 FILTER 수정자 사용

일부 경우 사용자는 서로 다른 차원에서 UV(고유 방문자) 수를 계산해야 할 수 있습니다. 예를 들어 Android UV, iPhone UV, Web UV, 전체 UV 말입니다. 많은 사용자가 이를 위해 CASE WHEN 을 선택합니다. 예:

SELECT
 day,
 COUNT(DISTINCT user_id) AS total_uv,
 COUNT(DISTINCT CASE WHEN flag IN ('android', 'iphone') THEN user_id ELSE NULL END) AS app_uv,
 COUNT(DISTINCT CASE WHEN flag IN ('wap', 'other') THEN user_id ELSE NULL END) AS web_uv
FROM T
GROUP BY day

그러나 이 경우 CASE WHEN 대신 FILTER 문법을 사용하는 것이 권장됩니다. FILTER 는 SQL 표준에 더 부합하고 훨씬 더 큰 성능 개선을 얻을 수 있기 때문입니다. FILTER 는 집계에 사용되는 값을 제한하기 위해 집계 함수에 적용되는 수정자입니다. 위 예제를 FILTER 수정자로 바꾸면 다음과 같습니다:

SELECT
 day,
 COUNT(DISTINCT user_id) AS total_uv,
 COUNT(DISTINCT user_id) FILTER (WHERE flag IN ('android', 'iphone')) AS app_uv,
 COUNT(DISTINCT user_id) FILTER (WHERE flag IN ('wap', 'other')) AS web_uv
FROM T
GROUP BY day

Flink SQL 최적화 도구는 같은 distinct 키에 대한 서로 다른 filter 인자를 인식할 수 있습니다. 예를 들어 위 예제에서 세 개의 COUNT DISTINCT 는 모두 user_id 열에 있습니다. 그러면 Flink 는 세 개의 상태 인스턴스 대신 공유 상태 인스턴스 하나를 사용해 상태 접근과 상태 크기를 줄일 수 있습니다. 일부 워크로드에서는 이로 인해 상당한 성능 개선을 얻을 수 있습니다.

MiniBatch 정규 조인 (MiniBatch Regular Joins)

기본적으로 정규 조인 operator 는 입력 레코드를 하나씩 처리합니다. 즉, (1) 현재 입력 레코드의 조인 키를 기반으로 상대편의 상태에서 관련 레코드를 조회하고, (2) 현재 입력 레코드를 추가하거나 재시도해 상태를 업데이트하고, (3) 현재 레코드와 관련 레코드에 따라 조인 결과를 출력합니다. 이 처리 패턴은 StateBackend 의 오버헤드를 증가시킬 수 있습니다(특히 RocksDB StateBackend 의 경우). 또한 특히 연쇄(cascading) 조인 시나리오에서 심각한 레코드 증폭을 초래해 너무 많은 중간 결과를 생성하고 성능 저하로 이어질 수 있습니다.

MiniBatch 조인은 위의 문제를 해결하려 합니다. 핵심 아이디어는 mini-batch 조인 operator 내부의 버퍼에 입력 묶음을 캐싱하는 것입니다. 버퍼가 지정된 크기 또는 시간 임계값에 도달하면 레코드가 조인 프로세스로 전달됩니다.

두 가지 핵심 최적화가 있습니다:

  • 조인 프로세스 전 데이터 수를 줄이기 위해 버퍼의 레코드를 접는(fold) 것.
  • 버퍼의 레코드가 처리될 때 중복 결과 출력을 최대한 억제하는 것.

예를 들어 다음 SQL 을 고려해 보세요:

SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5S';
SET 'table.exec.mini-batch.size' = '5000';

SELECT a.id as a_id, a.a_content, b.id as b_id, b.b_content
FROM a LEFT JOIN b
ON a.id = b.id

왼쪽과 오른쪽 입력 측 모두 조인 키 id 가 포함하는 고유 키를 가집니다(숫자는 id, 문자는 content 를 나타낸다고 가정).

MiniBatch 최적화는 정규 조인에 대해 기본적으로 비활성화되어 있습니다. 활성화하려면 table.exec.mini-batch.enabled, table.exec.mini-batch.allow-latency, table.exec.mini-batch.size 옵션을 설정해야 합니다. 자세한 내용은 configuration 페이지를 참고하세요.

여러 정규 조인 (Multiple Regular Joins)

Streaming

여러 비시간적(non-temporal) 정규 조인이 있는 스트리밍 Flink 작업은 큰 상태 크기로 인해 운영 불안정성과 성능 저하를 자주 겪습니다. 이는 종종 조인 체인에 의해 생성된 중간 상태가 입력 상태 자체보다 훨씬 크기 때문입니다. Flink 2.1 에서 레코드 증폭과 큰 중간 상태를 포함하는 조인 파이프라인을 위해 상태 크기를 크게 줄이고 성능을 개선하도록 설계된 새로운 다중 조인(multi-join) operator 를 도입했습니다. 이 새 operator 는 여러 입력 스트림에 걸쳐 조인을 동시에 처리함으로써 여러 테이블에 걸친 조인을 위해 중간 상태를 저장할 필요를 제거합니다. 이 "제로 중간 상태(zero intermediate state)" 접근은 주로 상태 축소를 목표로 하며, 일부 경우 자원 소비와 운영 안정성에서 상당한 이점을 제공합니다. 이 기법은 중간 상태를 필요 시 재평가하므로 저장 요구사항의 축소를 그에 상응하는 계산 노력의 증가와 맞바꿉니다.

대부분의 조인에서 처리 시간의 상당 부분은 상태에서 레코드를 가져오는 데 사용됩니다. MultiJoin operator 의 효율성은 이 중간 상태의 크기와 공통 조인 키의 선택성(selectivity)에 크게 의존합니다. 파이프라인이 레코드 증폭을 겪는 일반적인 시나리오—각 조인이 이전 것보다 더 많은 데이터와 레코드를 생성하는 경우—에서 MultiJoin operator 가 더 효율적입니다. 이는 operator 가 상호작용하는 상태를 훨씬 더 작게 유지해 더 안정적인 operator 로 이어지기 때문입니다. 조인 체인이 실제로 원래 레코드보다 적은 상태를 생성한다면 MultiJoin operator 는 여전히 전체적으로 더 적은 상태를 사용합니다. 그러나 이 특정 경우에는 최종 조인이 작동해야 하는 상태가 더 작으므로 이진 조인(binary join) 이 더 잘 수행될 수 있습니다.

MultiJoin Operator

MultiJoin operator 의 주요 이점:

  • 제로 중간 상태로 인해 상태 크기가 상당히 작아집니다.
  • 레코드 증폭이 있는 연쇄 조인의 성능이 개선됩니다.
  • 안정성 향상: 이진 조인의 다항식 증가 대신 처리된 레코드 수에 따라 선형 상태 증가.

또한 이진 조인 대신 MultiJoin 을 사용하는 파이프라인은 상태가 더 작고 노드가 더 적어 보통 초기화 및 복구 시간이 더 빠릅니다.

MultiJoin 을 언제 활성화할까?

작업에 공통 조인 키를 하나 이상 공유하는 여러 조인이 있고, 중간 조인의 중간 상태가 입력 소스보다 크다고 관찰된다면 MultiJoin operator 활성화를 고려하세요.

권장 사용 사례:

  • 공통 조인 키의 선택성이 높은 경우(키당 레코드 수가 적음)
  • 상당한 중간 상태가 있는 여러 연쇄 조인이 있는 문
  • 공통 조인 키에 상당한 데이터 왜곡이 없는 경우
  • 조인이 큰 상태를 생성하는 경우(상태 50+ GB)

공통 조인 키의 선택성이 낮으면(즉 같은 키 값을 공유하는 행 수가 많으면) MultiJoin operator 의 중간 상태 재계산이 성능에 심각한 영향을 줄 수 있습니다. 이러한 시나리오에서는 모든 조인 키를 사용해 데이터를 파티셔닝하는 이진 조인을 권장합니다.

MultiJoin 을 활성화하는 방법

적격한 모든 조인에 대해 이 최적화를 전역적으로 활성화하려면 다음 구성을 설정하세요:

SET 'table.optimizer.multi-join.enabled' = 'true';

또는 MULTI_JOIN 힌트를 사용해 특정 테이블에 대해 MultiJoin operator 를 활성화할 수 있습니다:

SELECT /*+ MULTI_JOIN(t1, t2, t3) */ * FROM t1
JOIN t2 ON t1.id = t2.id
JOIN t3 ON t1.id = t3.id;

힌트 접근 방식은 전역적으로 활성화하지 않고 특정 쿼리 블록에 MultiJoin 최적화를 선택적으로 적용할 수 있게 해줍니다. MULTI_JOIN 힌트에 대한 자세한 내용은 Join Hints 를 참고하세요. 구성 설정이 힌트보다 우선합니다.

중요: 현재 이것은 실험적 상태입니다 — 최적화와 호환성을 깨는 변경이 구현될 수 있습니다. 현재 스트리밍 INNER/LEFT 조인만 지원합니다. 레코드 파티셔닝으로 인해 조인 조건 사이에 공유되는 키가 하나 이상 필요합니다. 참고:

  • 지원: A JOIN B ON A.key = B.key JOIN C ON A.key = C.key (키로 파티셔닝)
  • 지원: A JOIN B ON A.key = B.key JOIN C ON B.key = C.key (추이성을 통한 키 파티셔닝)
  • 미지원: A JOIN B ON A.key1 = B.key1 JOIN C ON B.key2 = C.key2 (단일 operator 에서 A, B, C 를 함께 파티셔닝할 키가 없음. 이는 여러 MultiJoin operator 로 분할됨)

MultiJoin Operator 예제 - 벤치마크

다음은 기본 이진 조인과 MultiJoin operator 간의 10-way 벤치마크입니다. 첫 섹션에서 중간 상태의 양, 두 번째 섹션에서 operator 가 100% 바쁨에 도달할 때 처리된 레코드 수, 세 번째에서 체크포인트를 관찰할 수 있습니다.

위의 레코드 증폭을 포함하는 10-way 조인의 경우 상당한 개선을 관찰했습니다. 대략적인 수치는 다음과 같습니다:

  • 성능: 둘 다 100% 바쁨일 때 처리된 레코드가 2배에서 100배 이상 증가.
  • 상태 크기: 중간 상태가 커짐에 따라 3배에서 1000배 이상 감소.

MultiJoin operator 를 사용하면 전체 상태가 항상 더 작습니다. 이 경우 성능은 처음에는 동일하지만 중간 상태가 커짐에 따라 이진 조인의 성능은 저하되고 다중 조인은 안정적으로 유지되어 더 우수합니다.

이 10-way 조인에 대한 일반 벤치마크는 다음 구성으로 실행되었습니다: tenant_id 당 1 레코드(높은 선택성), 10 개의 upsert kafka 토픽, 병렬도 10, 토픽당 초당 1 레코드. unaligned 체크포인트와 incremental 체크포인트가 있는 rocksdb 를 사용했습니다. 각 작업은 8GB 프로세스 메모리, 1GB off-heap 메모리, 20% 네트워크 메모리를 포함하는 하나의 TaskManager 에서 실행되었습니다. JobManager 는 4GB 프로세스 메모리였습니다. 호스트 머신에는 M1 프로세서 칩, 32GB RAM, 1TB SSD 가 있었습니다. sink 는 조인만 벤치마크하도록 blackhole 커넥터를 사용했습니다. 벤치마크 데이터를 생성하는 데 사용된 SQL 은 다음과 같은 구조였습니다:

INSERT INTO JoinResultsMJ
SELECT *all fields*
FROM TenantKafka t
         LEFT JOIN SuppliersKafka s ON t.tenant_id = s.tenant_id AND ...
         LEFT JOIN ProductsKafka p ON t.tenant_id = p.tenant_id AND ...
         LEFT JOIN CategoriesKafka c ON t.tenant_id = c.tenant_id AND ...
         LEFT JOIN OrdersKafka o ON t.tenant_id = o.tenant_id AND ...
         LEFT JOIN CustomersKafka cust ON t.tenant_id = cust.tenant_id AND ...
         LEFT JOIN WarehousesKafka w ON t.tenant_id = w.tenant_id AND ...
         LEFT JOIN ShippingKafka sh ON t.tenant_id = sh.tenant_id AND ...
         LEFT JOIN PaymentKafka pay ON t.tenant_id = pay.tenant_id AND ...
         LEFT JOIN InventoryKafka i ON t.tenant_id = i.tenant_id AND ...;

Delta 조인 (Delta Joins)

스트리밍 작업에서 정규 조인은 정확성을 보장하기 위해 두 입력의 모든 과거 데이터를 유지합니다. 시간이 지남에 따라 이로 인해 상태가 계속 커지고, 자원 사용량이 증가하며 안정성에 영향을 줍니다.

이러한 문제를 완화하기 위해 Flink 는 delta 조인 operator 를 도입합니다. 핵심 아이디어는 정규 조인이 유지하는 큰 상태를 소스 테이블의 데이터를 직접 재사용하는 양방향 조회 기반 조인으로 대체하는 것입니다. 전통적인 정규 조인과 비교해 delta 조인은 상태 크기를 상당히 줄이고, 작업 안정성을 향상시키며, 전체 자원 소비를 낮춥니다.

이 기능은 기본적으로 활성화되어 있습니다. 다음 조건이 모두 충족되면 정규 조인이 자동으로 delta 조인으로 최적화됩니다:

  • SQL 패턴이 최적화 기준을 충족하는 경우. 자세한 내용은 지원 기능 및 제한 사항 을 참고하세요.
  • 소스 테이블의 외부 스토리지 시스템이 delta 조인용 빠른 쿼리를 위한 인덱스 정보를 제공하는 경우. 현재 Apache Fluss (Incubating) 가 Flink 에 테이블 수준의 인덱스 정보를 제공하여 이러한 테이블을 delta 조인의 소스 테이블로 사용할 수 있게 합니다. 자세한 내용은 Fluss 문서 를 참고하세요.

작동 원리

Flink 에서 정규 조인은 두 입력 측의 모든 수신 레코드를 상태에 저장해, 반대편에서 데이터가 도착할 때 해당 레코드를 올바르게 매칭할 수 있도록 보장합니다.

반면 delta 조인은 외부 스토리지 시스템의 인덱싱 기능을 활용합니다. 상태 조회를 수행하는 대신 delta 조인은 외부 스토리지에 대해 효율적인 인덱스 기반 쿼리를 직접 발행해 일치하는 레코드를 검색합니다. 이 접근 방식은 Flink 상태와 외부 시스템 간의 중복 데이터 저장을 제거합니다.

중요한 구성 (Important Configurations)

Delta 조인 최적화는 기본적으로 활성화되어 있습니다. 다음 구성을 설정해 이 기능을 수동으로 비활성화할 수 있습니다:

SET 'table.optimizer.delta-join.strategy' = 'NONE';

자세한 내용은 Configuration 페이지를 참고하세요.

Delta 조인 성능을 세밀하게 조정하려면 다음 매개변수도 구성할 수 있습니다:

  • table.exec.delta-join.cache-enabled
  • table.exec.delta-join.left.cache-size
  • table.exec.delta-join.right.cache-size

자세한 내용은 Configuration 페이지를 참고하세요.

지원 기능 및 제한 사항

Delta 조인은 지속적으로 진화하고 있으며 현재 다음 기능을 지원합니다.

  • INSERT-only 테이블을 소스 테이블로 지원.
  • DELETE 연산이 없는 CDC 테이블을 소스 테이블로 지원.
  • 소스와 delta 조인 사이의 projectionfilter 연산 지원.
  • delta 조인 operator 내부의 caching 지원.
  • 연쇄 delta 조인(cascaded delta joins) 지원 — 쿼리의 적격한 조인 노드가 소스에서 sink 로 순차적으로 delta 조인 노드로 변환됨.
  • delta 조인 이후의 lookup join 지원.
  • delta 조인과 다운스트림 operator 사이의 projection 및 filter 에서 비결정적 함수 지원.

그러나 Delta 조인에도 여러 제한 사항 이 있습니다. 다음 조건 중 하나라도 포함하는 작업은 delta 조인으로 최적화될 수 없습니다:

  • 테이블의 인덱스 키 가 조인의 동치 조건(equivalence conditions) 에 포함되어야 합니다.
  • 현재 INNER JOIN 만 지원됩니다.
  • 다운스트림 operator 가 중복 변경(duplicate changes) 을 처리할 수 있어야 합니다. 예: upsertMaterialize 없이 UPSERT 모드 로 동작하는 sink.
  • CDC 스트림 을 소비할 때 조인 키upsert 키[^1]의 일부여야 합니다.
  • CDC 스트림 을 소비할 때 모든 filterupsert 키[^1]에 적용되어야 합니다.
  • 소스와 delta 조인 사이의 projection 또는 filter 에서 비결정적 함수 는 허용되지 않습니다.

[^1]: Flink 는 upsert 키를 보강하기 위해 테이블 수준의 immutable columns(불변 열) 제약 조건을 정의하는 것을 지원해 더 많은 시나리오에서 delta 조인 최적화를 가능하게 합니다. immutable columns 제약 조건은 주어진 기본 키에 대해 한 번 설정되면 수정될 수 없는 특정 열을 선언합니다. immutable columns 는 기본 키와 결합되어 다운스트림 operator 로 전파되는 새 upsert 키를 형성합니다. 이 정보는 외부 스토리지 시스템이 제공합니다. Apache Fluss (Incubating) 는 향후 테이블 수준의 immutable columns 제약 조건을 지원할 계획입니다.

더 알아보기 (Learn more)