MERGE 지원

MERGE 지원 (Supporting MERGE)

Trino 엔진은 행 수준 SQL MERGE를 지원하는 API를 제공해요. MERGE를 구현하려면 커넥터가 다음을 제공해야 해요.

  • 보통 ConnectorPageSink 위에 계층화되는 ConnectorMergeSink 구현.
  • "rowId" 컬럼 핸들을 얻고, 행 변경 패러다임을 얻고, MERGE 연산을 시작·완료하는 ConnectorMetadata의 메서드들.

출처: 문서

본문

SQL MERGE를 구현하는 데 쓰는 Trino 엔진 메커니즘은 SQL DELETEUPDATE를 지원하는 데도 사용돼요. 즉 커넥터가 SQL MERGE 지원만 구현하면 모든 데이터 조작 언어(DML) 연산을 얻을 수 있어요.

표준 SQL MERGE (Standard SQL MERGE)

쿼리 엔진마다 SQL MERGE의 정의가 다를 수 있어요. Trino는 2016년에 발표된 엄격한 SQL 명세 ISO/IEC 9075를 지원해요. 간단한 예로, target_tablesource_table 테이블이 다음과 같이 정의돼 있다고 해요.

CREATE TABLE accounts (
    customer VARCHAR,
    purchases DECIMAL,
    address VARCHAR);
INSERT INTO accounts (customer, purchases, address) VALUES ...;
CREATE TABLE monthly_accounts_update (
    customer VARCHAR,
    purchases DECIMAL,
    address VARCHAR);
INSERT INTO monthly_accounts_update (customer, purchases, address) VALUES ...;

monthly_accounts_update에서 accounts로의 가능한 MERGE 연산은 다음과 같아요.

MERGE INTO accounts t USING monthly_accounts_update s
    ON (t.customer = s.customer)
    WHEN MATCHED AND s.address = 'Berkeley' THEN
        DELETE
    WHEN MATCHED AND s.customer = 'Joe Shmoe' THEN
        UPDATE SET purchases = purchases + 100.0
    WHEN MATCHED THEN
        UPDATE
            SET purchases = s.purchases + t.purchases, address = s.address
    WHEN NOT MATCHED THEN
        INSERT (customer, purchases, address)
            VALUES (s.customer, s.purchases, s.address);

SQL MERGE는 각 WHEN 절을 소스 순서로 매칭하려 해요. 매치가 발견되면 해당 DELETE, INSERT, UPDATE가 실행되고 이후 WHEN 절은 무시돼요.

소스 테이블이나 쿼리의 행이 대상 테이블의 행과 일치하면 SQL MERGE는 대상 테이블과 소스에 대해 두 가지 연산을 지원해요.

  • UPDATE: 대상 행의 컬럼을 갱신.
  • DELETE: 대상 행을 삭제.

NOT MATCHED 경우에는 SQL MERGEINSERT 연산만 지원해요. 삽입되는 값은 임의적이지만 보통 소스 테이블이나 쿼리의 매치되지 않은 행에서 와요.

RowChangeParadigm

커넥터마다 기반 저장 시스템이 부과하는 행 갱신 표현 방식이 달라요. Trino 엔진은 이런 다른 패러다임을 ConnectorMetadata.getRowChangeParadigm(...) 메서드가 반환하는 RowChangeParadigm 열거형의 요소로 분류해요.

RowChangeParadigm 열거형 값은 다음과 같아요.

  • CHANGE_ONLY_UPDATED_COLUMNS: rowId로 식별된 행의 개별 컬럼을 갱신할 수 있는 커넥터용. 대응하는 병합 프로세서 클래스는 ChangeOnlyUpdatedColumnsMergeProcessor.
  • DELETE_ROW_AND_INSERT_ROW: 행 변경을 행 삭제와 행 삽입의 쌍으로 나타내는 커넥터용. 대응하는 병합 프로세서 클래스는 DeleteAndInsertMergeProcessor.

MERGE 처리 개요 (Overview of MERGE processing)

MERGE 문은 대상 테이블과 소스 사이에 MERGE 조건으로 RIGHT JOIN을 만들어 처리해요. 소스는 테이블일 수도 임의의 쿼리일 수도 있어요. 소스 테이블이나 쿼리의 각 행에 대해 MERGE는 다음을 담은 ROW 객체를 만들어요.

  • UPDATEINSERT 경우의 데이터 컬럼 값. DELETE 경우에는 파티셔닝과 버케팅을 결정하는 파티션 컬럼만 null이 아니에요.
  • 일부 대상 행과 매치된 소스 행이면 true, 그렇지 않으면 false인 boolean 컬럼.
  • 병합 경우 연산이 UPDATE, DELETE, INSERT인지, 아니면 매치된 경우가 없는 소스 행인지 식별하는 정수. 소스 행이 어떤 병합 경우와도 매치되지 않으면 분배를 결정하는 컬럼 외의 모든 데이터 컬럼 값은 null이고 연산 번호는 -1이에요.

RIGHT JOIN 결과에서 SearchedCaseExpression을 구성해 MERGEWHEN 절을 나타내요. 위 예시는 MERGE가 마치 SearchedCaseExpression이 다음과 같이 쓰인 것처럼 실행돼요.

SELECT
 CASE
   WHEN present AND s.address = 'Berkeley' THEN
       -- Null values for delete; present=true; operation DELETE=2, case_number=0
       row(null, null, null, true, 2, 0)
   WHEN present AND s.customer = 'Joe Shmoe' THEN
       -- Update column values; present=true; operation UPDATE=3, case_number=1
       row(t.customer, t.purchases + 100.0, t.address, true, 3, 1)
   WHEN present THEN
       -- Update column values; present=true; operation UPDATE=3, case_number=2
       row(t.customer, s.purchases + t.purchases, s.address, true, 3, 2)
   WHEN (present IS NULL) THEN
       -- Insert column values; present=false; operation INSERT=1, case_number=3
       row(s.customer, s.purchases, s.address, false, 1, 3)
   ELSE
       -- Null values for no case matched; present=false; operation=-1,
       --     case_number=-1
       row(null, null, null, false, -1, -1)
 END
 FROM (SELECT *, true AS present FROM target_table) t
   RIGHT JOIN source_table s ON s.customer = t.customer;

Trino 엔진은 RIGHT JOINCASE 표현식을 실행하고, 대상 테이블 행이 하나 이상의 소스 표현식 행과 일치하지 않도록 보장하며, 최종적으로 ConnectorMergeSink.storeMergedRows(...) 메서드를 실행하는 노드로 라우팅될 페이지 시퀀스를 만들어요.

DELETEUPDATE처럼 MERGE 대상 테이블 행은 커넥터별 rowId 컬럼 핸들로 식별돼요. MERGE의 경우 rowId 핸들은 ConnectorMetadata.getMergeRowIdColumnHandle(...)가 반환해요.

MERGE 재분배 (MERGE redistribution)

Trino MERGE 구현은 UPDATE가 파티셔닝·버케팅을 결정하는 컬럼의 값을 바꿀 수 있게 해주므로, 병합된 파티셔닝·버케팅 컬럼으로 행을 쓰는 책임을 가진 워커 노드로 MERGE 연산의 행을 "재분배"해야 해요.

MERGE 과정은 일반적으로 Trino 노드 사이에서 병합된 행의 재분배가 필요하므로, 저장될 페이지의 행 순서는 불확정적이에요. 삭제된 행에 대해 오름차순 rowId 순서에 의존하는 Hive 같은 커넥터는 저장 전에 삭제된 행을 정렬해야 해요.

주어진 파티션의 모든 삽입 행이 단일 노드에 도달하도록, 파티션 키/버킷 컬럼에 대한 재분배 해시가 페이지 파티션 키에 적용돼요. 해시 결과로 특정 파티션/버킷의 모든 행은 MATCHED 행이든 NOT MATCHED 행이든 함께 해시돼요.

RowChangeParadigmDELETE_ROW_AND_INSERT_ROW인 커넥터의 경우 삽입된 행은 ConnectorMetadata.getInsertLayout()가 제공하는 레이아웃으로 분배돼요. 일부 커넥터는 갱신된 행에도 같은 레이아웃을 사용해요. 다른 커넥터는 ConnectorMetadata.getUpdateLayout()이 제공하는 갱신된 행용 특수 레이아웃을 필요로 해요.

커넥터의 MERGE 지원 (Connector support for MERGE)

MERGE 처리를 시작하려면 Trino 엔진이 다음을 호출해요.

  • ConnectorMetadata.getMergeRowIdColumnHandle(...): rowId 컬럼 핸들 획득.
  • ConnectorMetadata.getRowChangeParadigm(...): 기존 테이블 행 변경에 대해 커넥터가 지원하는 패러다임 획득.
  • ConnectorMetadata.beginMerge(...): 병합 연산용 ConnectorMergeTableHandle 획득. 이 ConnectorMergeTableHandle 객체는 커넥터가 MERGE 연산을 지정하는 데 필요한 모든 정보를 담아요.
  • ConnectorMetadata.getInsertLayout(...): 여기서 쓰기 재분배에 영향을 주는 파티션 또는 테이블 컬럼 목록을 추출.
  • ConnectorMetadata.getUpdateLayout(...): 이 레이아웃이 비어 있지 않으면 MERGE 연산에서 나온 갱신된 행을 분배하는 데 사용.

해시의 대상인 노드들에서 Trino 엔진은 ConnectorPageSinkProvider.createMergeSink(...)을 호출해 ConnectorMergeSink를 만들어요.

병합된 행의 각 페이지를 쓰려면 Trino 엔진이 ConnectorMergeSink.storeMergedRows(Page)를 호출해요. storeMergedRows(Page) 메서드는 페이지의 행을 반복하며 MATCHED 경우에는 갱신·삭제를, NOT MATCHED 경우에는 삽입을 수행해요.

RowChangeParadigm.DELETE_ROW_AND_INSERT_ROW를 사용할 때 엔진은 storeMergedRows(Page)가 호출되기 전에 UPDATE 연산을 DELETEINSERT 연산의 쌍으로 변환해요.

MERGE 연산을 완료하려면 Trino 엔진이 ConnectorMetadata.finishMerge(...)를 호출하며 테이블 핸들과 Slice 인스턴스로 인코딩된 JSON 객체 컬렉션을 전달해요. 이 객체들은 MERGE 연산으로 무엇이 바뀌었는지 지정하는 커넥터별 정보를 담아요. 보통 이 JSON 객체는 MERGE 연산이 만든 파일과 테이블·파티션 통계를 포함해요. 커넥터는 있다면 적절한 조치를 취해요.

MERGE용 RowChangeProcessor 구현 (RowChangeProcessor implementation for MERGE)

MERGE 구현에서 각 RowChangeParadigmRowChangeProcessor 인터페이스를 구현하는 내부 Trino 엔진 클래스에 대응해요. RowChangeProcessor에는 흥미로운 메서드가 하나 있어요: Page transformPage(Page). 출력 페이지의 형식은 RowChangeParadigm에 따라 달라요.

커넥터는 RowChangeProcessor 인스턴스에 접근할 수 없어요. 이는 커넥터가 선택한 RowChangeParadigm에 따라 병합 페이지 행을 저장할 행으로 변환하기 위해 Trino 엔진 내부에서 사용돼요.

transformPage()에 제공된 페이지는 다음으로 구성돼요.

  • 있다면 쓰기 재분배 컬럼
  • 파티션되거나 버킷된 테이블의 경우, long 해시 값 컬럼
  • 매치된 경우 대상 테이블의 행에 대한 rowId 컬럼, 매치되지 않으면 null
  • 병합 경우 RowBlock
  • 정수 경우 번호 블록
  • 구분되지 않으면 값 0인 바이트 is_distinct 블록

병합 경우 RowBlock은 다음 레이아웃을 가져요.

  • 파티션 컬럼을 포함한 테이블의 각 컬럼에 대한 블록. 테이블 컬럼 순서로.
  • 소스 행이 대상 행과 매치되었으면 true, 아니면 false인 boolean "present" 값을 담은 블록.
  • INSERT = 1, DELETE = 2, UPDATE = 3으로 인코딩된 MERGE 경우 연산 번호를 담은 블록. 매치된 MERGE 경우가 없으면 -1.
  • 행에 대해 매치된 WHEN 절의, 0부터 시작하는 MERGE 경우 번호를 담은 블록. 매치된 절이 없으면 -1.

transformPage가 반환한 페이지는 다음으로 구성돼요.

  • 테이블 컬럼 순서의 모든 테이블 컬럼.
  • tinyint 타입 병합 경우 연산 블록.
  • 정수 타입 병합 경우 번호 블록.
  • rowId 블록은 제공된 입력 페이지에서 변경되지 않음.
  • 행이 업데이트 연산에서 파생된 삽입이면 1, 아니면 0인 바이트 블록. 이 블록은 갱신과 삭제 그리고 삽입을 나타내는 커넥터에서 변경된 행 수를 올바르게 계산하는 데 사용돼요.

transformPage는 반환하는 페이지에 연산 번호가 -1인 행이 없도록 보장해야 해요.

중복 일치 대상 행 감지 (Detecting duplicate matching target rows)

SQL MERGE 명세는 각 MERGE 경우에서 단일 대상 테이블 행이 MERGE 경우 조건 표현식을 적용한 뒤 최대 하나의 소스 행과 일치해야 한다고 요구해요. 이 오류를 찾는 첫 단계는 대상 테이블 스캔 위의 AssignUniqueId 노드를 사용해 대상 테이블의 각 행에 고유 id를 붙이는 거예요. RIGHT JOIN의 투영 결과에는 매치된 대상 테이블 행의 고유 id와 WHEN 절 번호가 있어요. MarkDistinct 노드는 다른 행이 같은 고유 id와 WHEN 절 번호를 가지지 않으면 true, 아니면 false인 is_distinct 컬럼을 추가해요. 어떤 행의 is_distinct가 false와 같으면 MERGE_TARGET_ROW_MULTIPLE_MATCHES 예외가 발생하고 MERGE 연산이 실패해요.

ConnectorMergeTableHandle API

ConnectorMergeTableHandle 인터페이스는 ConnectorMetadata.beginMerge()에 원래 전달된 ConnectorTableHandle을 검색하는 getTableHandle() 메서드 하나를 정의해요.

ConnectorPageSinkProvider API

SQL MERGE를 지원하려면 ConnectorPageSinkProviderConnectorMergeSink를 만드는 메서드를 구현해야 해요.

createMergeSink:

ConnectorMergeSink createMergeSink(
    ConnectorTransactionHandle transactionHandle,
    ConnectorSession session,
    ConnectorMergeTableHandle mergeHandle)

ConnectorMergeSink API

MERGE를 지원하려면 커넥터가 보통 커넥터의 ConnectorPageSink 위에 계층화된 ConnectorMergeSink 구현을 정의해야 해요.

ConnectorMergeSinkConnectorPageSinkProvider.createMergeSink() 호출로 만들어져요.

흥미로운 메서드는 다음과 같아요.

storeMergedRows:

void storeMergedRows(Page page)

Trino 엔진은 ConnectorPageSinkProvider.createMergeSink()이 반환한 ConnectorMergeSink 인스턴스의 storeMergedRows(Page) 메서드를 호출하며 RowChangeProcessor.transformPage() 메서드가 만든 페이지를 전달해요. 그 페이지는 테이블 컬럼 순서의 모든 테이블 컬럼, 그 다음 TINYINT 연산 컬럼, 그 다음 INTEGER 병합 경우 번호 컬럼, 그 다음 rowId 컬럼으로 구성돼요.

storeMergedRows()의 역할은 페이지의 행을 반복하고, 연산 컬럼 값(INSERT, DELETE, UPDATE)에 따라 처리하거나 행을 무시하는 거예요. 적절한 패러다임을 선택하면 커넥터는 UPDATE 연산을 DELETEINSERT 연산으로 변환되도록 요청할 수 있어요.

finish:

CompletableFuture<Collection<Slice>> finish()

Trino 엔진은 특정 ConnectorMergeSink 인스턴스가 모든 데이터를 처리하면 finish()를 호출해요. 커넥터는 처리된 행에 대한 커넥터별 정보를 나타내는 Slice 컬렉션을 담은 future를 반환해요. 보통 이는 행 수를 포함하고, 생성되거나 변경된 파일·파티션 같은 정보를 포함할 수도 있어요.

ConnectorMetadata MERGE API

MERGE를 구현하는 커넥터는 이 ConnectorMetadata 메서드들을 구현해야 해요.

getRowChangeParadigm():

RowChangeParadigm getRowChangeParadigm(
    ConnectorSession session,
    ConnectorTableHandle tableHandle)

이 메서드는 엔진이 MERGE 문 처리를 시작할 때 호출돼요. 커넥터는 RowChangeParadigm 열거형 인스턴스를 반환해야 해요. 커넥터가 MERGE를 지원하지 않는다면, 커넥터가 SQL MERGE를 지원하지 않는다는 NOT_SUPPORTED 예외를 던져야 해요. 기본 구현은 메서드가 구현되지 않았을 때 이미 이 예외를 던진다는 점에 유의하세요.

getMergeRowIdColumnHandle():

ColumnHandle getMergeRowIdColumnHandle(
    ConnectorSession session,
    ConnectorTableHandle tableHandle)

이 메서드는 MERGE 문의 쿼리 계획 초기 단계에서 호출돼요. 반환된 ColumnHandle은 커넥터가 병합할 행을 식별하는 데 사용하는 rowId와, 커넥터가 MERGE 연산을 완료하는 데 필요한 행의 다른 필드를 제공해요.

getInsertLayout():

Optional<ConnectorTableLayout> getInsertLayout(
    ConnectorSession session,
    ConnectorTableHandle tableHandle)

이 메서드는 쿼리 계획 중에 MERGE 연산이 삽입하는 행에 사용할 테이블 레이아웃을 얻기 위해 호출돼요. 일부 커넥터에서는 이 레이아웃이 삭제된 행에도 사용돼요.

getUpdateLayout():

Optional<ConnectorTableLayout> getUpdateLayout(
    ConnectorSession session,
    ConnectorTableHandle tableHandle)

이 메서드는 쿼리 계획 중에 MERGE 연산이 삭제하는 행에 사용할 테이블 레이아웃을 얻기 위해 호출돼요. 선택적 반환 값이 있으면 Trino 엔진은 갱신된 행에 이 레이아웃을 사용해요. 그렇지 않으면 ConnectorMetadata.getInsertLayout의 결과를 사용해 갱신된 행을 분배해요.

beginMerge():

ConnectorMergeTableHandle beginMerge(
     ConnectorSession session,
     ConnectorTableHandle tableHandle)

MERGE 실행 계획 만들기의 마지막 단계로 커넥터의 beginMerge() 메서드가 호출되며 sessiontableHandle이 전달돼요.

beginMerge()는 커넥터에서 MERGE 처리를 시작하기 위해 필요한 모든 오케스트레이션을 수행해요. 이 오케스트레이션은 커넥터마다 달라요. 예를 들어 트랜잭션 테이블에서 동작하는 Hive 커넥터의 경우 beginMerge()는 테이블이 트랜잭션인지 확인하고 Hive Metastore 트랜잭션을 시작해요.

beginMerge()는 핸들이 finishMerge()와 스플릿 생성 메커니즘으로 전달될 때 커넥터가 필요한 추가 정보와 함께 ConnectorMergeTableHandle을 반환해요. 대부분의 커넥터에서 반환된 테이블 핸들은 적어도 이 테이블 핸들이 MERGE 연산용 테이블 핸들임을 식별하는 플래그를 포함해요.

finishMerge():

void finishMerge(
    ConnectorSession session,
    ConnectorMergeTableHandle tableHandle,
    Collection<Slice> fragments)

MERGE 처리 중 Trino 엔진은 ConnectorMergeSink.finish()가 반환한 Slice 컬렉션을 축적해요. 엔진은 finishMerge()를 호출하며 테이블 핸들과 그 Slice 조각 컬렉션을 전달해요. 응답으로 커넥터는 MERGE 연산을 완료하기 위한 적절한 조치를 취해요. 이런 조치에는 있다면 기반 트랜잭션 커밋이나 다른 리소스 해제가 포함될 수 있어요.

더 알아보기 (Learn more)

MERGE는 트랜잭션 커넥터에서 특히 중요해요. Hive 커넥터와 같은 트랜잭션 커넥터에서 MERGE가 어떻게 동작하는지 함께 살펴보세요.