커넥터

커넥터 (Connectors)

커넥터는 Trino에서 쿼리의 모든 데이터 원천이에요. 데이터 소스에 이를 뒷받침하는 테이블이 없더라도, Trino가 기대하는 API에 맞추기만 하면 그 데이터에 대해 쿼리를 작성할 수 있어요.

출처: 문서

본문

ConnectorFactory

커넥터의 인스턴스는 Trino가 플러그인에서 getConnectorFactory()를 호출할 때 생성되는 ConnectorFactory 인스턴스가 만들어요. 커넥터 팩토리는 커넥터 이름을 제공하고 Connector 객체 인스턴스를 만드는 책임을 가진 단순한 인터페이스예요. 읽기만 지원하고 쓰기는 지원하지 않는 기본 커넥터 구현은 다음 서비스의 인스턴스를 반환해야 해요.

  • ConnectorMetadata
  • ConnectorSplitManager
  • ConnectorRecordSetProvider 또는 ConnectorPageSourceProvider

설정 (Configuration)

커넥터 팩토리의 create() 메서드는 카탈로그 속성 파일의 모든 속성을 담은 config 맵을 받아요. 이 맵은 커넥터를 설정하는 데 쓸 수 있지만, 모든 값이 문자열이므로 다른 데이터 타입을 나타낸다면 추가 처리가 필요할 수 있어요. 또한 제공된 모든 속성이 알려진 속성인지 검증하지 않아요. 사용자가 속성 이름을 잘못 입력해 커넥터가 속성을 무시하면 커넥터가 예상과 다르게 동작할 수 있어요.

설정을 더 견고하게 만들려면 Configuration 클래스를 정의하세요. 이 클래스는 사용 가능한 모든 속성, 그 타입, 추가 검증 규칙을 설명해요.

import io.airlift.configuration.Config;
import io.airlift.configuration.ConfigDescription;
import io.airlift.configuration.ConfigSecuritySensitive;
import io.airlift.units.Duration;
import io.airlift.units.MaxDuration;
import io.airlift.units.MinDuration;

import javax.validation.constraints.NotNull;

public class ExampleConfig
{
    private String secret;
    private Duration timeout = Duration.succinctDuration(10, TimeUnit.SECONDS);

    public String getSecret()
    {
        return secret;
    }

    @Config("secret")
    @ConfigDescription("Secret required to access the data source")
    @ConfigSecuritySensitive
    public ExampleConfig setSecret(String secret)
    {
        this.secret = secret;
        return this;
    }

    @NotNull
    @MaxDuration("10m")
    @MinDuration("1ms")
    public Duration getTimeout()
    {
        return timeout;
    }

    @Config("timeout")
    public ExampleConfig setTimeout(Duration timeout)
    {
        this.timeout = timeout;
        return this;
    }
}

위 예시는 두 설정 속성을 정의하고 다음 방법으로 커넥터를 더 견고하게 만들어요.

  • 지원되는 모든 속성을 정의해 서버 시작 시 설정의 철자 오류를 감지
  • 기본 타임아웃 값을 정의해 연결이 무한정 멈추지 않게 방지
  • 0ms처럼 모든 요청을 실패하게 만드는 잘못된 타임아웃 값을 방지
  • 서로 다른 단위의 타임아웃 값을 파싱하고 잘못된 값을 감지
  • 시크릿 값을 평문으로 로깅하지 못하게 방지

설정 클래스는 Guice 모듈에 바인딩되어야 해요.

import com.google.inject.Binder;
import com.google.inject.Module;

import static io.airlift.configuration.ConfigBinder.configBinder;

public class ExampleModule
        implements Module
{
    public ExampleModule()
    {
    }

    @Override
    public void configure(Binder binder)
    {
        configBinder(binder).bindConfig(ExampleConfig.class);
    }
}

그리고 커넥터의 새 인스턴스를 만들 때 커넥터 팩토리에서 모듈을 초기화해야 해요.

@Override
public Connector create(String catalogName, Map<String, String> config, ConnectorContext context)
{
    requireNonNull(config, "config is null");
    Bootstrap app = new Bootstrap("io.trino.bootstrap.catalog." + catalogName, new ExampleModule());
    Injector injector = app
            .doNotInitializeLogging()
            .setRequiredConfigurationProperties(config)
            .initialize();

    return injector.getInstance(ExampleConnector.class);
}

참고

카탈로그 속성 파일의 환경 변수(예: secret=${ENV:SECRET})는 io.airlift.bootstrap.Bootstrap 클래스로 모듈을 초기화할 때만 해석돼요. 자세한 내용은 시크릿 문서를 참조하세요.

한 속성만 바꾸기 위해 같은 커넥터로 여러 카탈로그를 정의해야 한다면, 스키마·테이블 속성 지원 추가를 고려해 보세요. 그렇게 하면 더 세밀한 설정이 가능해요. 커넥터가 스키마 관리를 지원하지 않는다면, 선택한 컬럼의 쿼리 프레디킷을 런타임에 필요한 설정을 전달하는 방법으로 쓸 수 있어요.

예를 들어 Git 리포지토리의 커밋을 읽는 커넥터를 만들 때, 리포지토리 URL은 설정 속성일 수 있어요. 하지만 그러면 카탈로그는 단일 리포지토리의 데이터만 반환할 수 있게 돼요. 대안으로 URL을 컬럼으로 두면 모든 select 쿼리에 그에 대한 프레디킷이 필요해요.

SELECT *
FROM git.default.commits
WHERE url = 'https://github.com/trinodb/trino.git'

ConnectorMetadata

커넥터 메타데이터 인터페이스는 Trino가 특정 데이터 소스에 대한 스키마, 테이블, 컬럼 및 기타 메타데이터 목록을 얻을 수 있게 해줘요.

기본 읽기 전용 커넥터는 다음 메서드를 구현해야 해요.

  • listSchemaNames
  • listTables
  • streamTableColumns
  • getTableHandle
  • getTableMetadata
  • getColumnHandles
  • getColumnMetadata

더 많은 메서드를 구현하는 전략에 관심이 있다면 예제 HTTP 커넥터와 Cassandra 커넥터를 살펴보세요. 기반 데이터 소스가 스키마, 테이블, 컬럼을 지원한다면 이 인터페이스는 구현하기 쉬울 거예요. 예제 HTTP 커넥터처럼 관계형 데이터베이스가 아닌 것을 적응시키려 한다면, 데이터 소스를 Trino의 스키마·테이블·컬럼 개념에 매핑하는 방법을 창의적으로 고민해야 할 수 있어요.

커넥터 메타데이터 인터페이스는 다른 커넥터 기능도 구현할 수 있게 해줘요.

  • 스키마 관리. 스키마, 테이블, 테이블 컬럼, 뷰, 구체화된 뷰를 생성·변경·삭제하는 기능.

  • 테이블·컬럼 코멘트와 속성 지원.

  • 스키마, 테이블, 뷰 인가(authorization).

  • 테이블 함수 실행.

  • 비용 기반 옵티마이저(CBO)가 쓰는 테이블 통계 제공과 쓰기 중 및 선택된 테이블 분석 시 통계 수집.

  • 데이터 수정. 즉 테이블 행 삽입·업데이트·삭제, 구체화된 뷰 리프레시, 전체 테이블 잘라내기, 쿼리 결과로 테이블 만들기.

  • 역할(role)과 그랜트 관리.

  • 푸시다운:

    • Limit와 Top N - 정렬 항목이 있는 limit
    • 프레디킷
    • 프로젝션
    • 샘플링
    • 집계
    • 조인
    • 테이블 함수 호출

데이터 수정에는 ConnectorPageSinkProvider 구현도 필요하다는 점에 유의하세요.

Trino가 SELECT 쿼리를 받으면 이를 중간 표현(IR, Intermediate Representation)으로 파싱해요. 그런 다음 최적화 동안 SQL 절과 관련된 연산을 커넥터가 처리할 수 있는지 ConnectorMetadata 서비스의 다음 메서드 중 하나를 호출해 확인해요.

  • applyLimit
  • applyTopN
  • applyFilter
  • applyProjection
  • applySample
  • applyAggregation
  • applyJoin
  • applyTableFunction
  • applyTableScanRedirect

커넥터는 특정 푸시다운을 지원하지 않거나 동작이 효과가 없었다는 것을 Optional.empty()를 반환해 나타낼 수 있어요. 커넥터는 주어진 쿼리 최적화 중에 이 메서드들이 여러 번 호출될 수 있다는 것을 예상해야 해요.

경고

이 메서드 호출이 그 호출에 대해 효과가 없으면, 커넥터가 일반적으로 특정 푸시다운을 지원하더라도 Optional.empty()를 반환하는 것이 중요해요. 그렇지 않으면 옵티마이저가 무한 루프에 빠질 수 있어요.

그 외에는 이 메서드들이 새 테이블 핸들을 담은 결과 객체를 반환해요. 새 테이블 핸들은 테이블 스캔 노드가 만든 테이블에 연산(filter, project, limit 등)을 적용해 파생된 가상 테이블을 나타내요. 쿼리가 실제로 실행되면 ConnectorRecordSetProvider 또는 ConnectorPageSourceProviderConnectorTableHandle로 푸시다운된 모든 최적화를 사용할 수 있어요.

반환된 테이블 핸들은 이후 커넥터가 구현하는 ConnectorRecordSetProviderConnectorPageSourceProvider 같은 다른 서비스로 전달돼요.

Limit와 top-N 푸시다운 (Limit and top-N pushdown)

LIMITORDER BY 절이 있는 SELECT 쿼리를 실행할 때 쿼리 계획에는 SortLimit 연산이 포함될 수 있어요.

계획에 SortLimit 연산이 모두 있으면 엔진은 커넥터 메타데이터 서비스의 applyTopN 메서드를 호출해 limit을 커넥터로 푸시다운하려고 해요. Sort 연산이 없고 Limit만 있으면 applyLimit 메서드가 호출되고, 커넥터는 임의의 순서로 결과를 반환할 수 있어요.

이 메서드들에 전달된 정보가 유용하지만 제공된 limit보다 적은 행을 만들 수 있다는 것이 보장되지 않는다면, 파생 테이블에 대한 새 핸들과 limitGuaranteed(LimitApplicationResult 안) 또는 topNGuaranteed(TopNApplicationResult 안) 플래그를 false로 설정한 비어 있지 않은 결과를 반환해야 해요.

커넥터가 제공된 limit보다 적은 행을 만든다는 것을 보장할 수 있다면 "limit guaranteed" 또는 "topN guaranteed" 플래그를 true로 설정한 비어 있지 않은 결과를 반환해야 해요.

참고

applyTopNSort 연산에서 정렬 항목을 받는 유일한 메서드예요.

쿼리에서 ORDER BY 부분은 어떤 컬럼이든 어떤 순서로든 포함할 수 있어요. 하지만 커넥터의 데이터 소스는 제한된 조합만 지원할 수 있어요. 플러그인 작성자는 커넥터가 푸시다운을 무시하고 모든 데이터를 반환해 엔진이 정렬하게 할지, 아니면 특정 순서가 지원되지 않는다는 것을 사용자에게 알리기 위해 예외를 던질지 결정해야 해요(특히 모든 데이터를 가져오는 것이 너무 비싸거나 시간이 오래 걸린다면). 예외를 던질 때는 TrinoException 클래스를 INVALID_ORDER_BY 오류 코드와 실행 가능한 메시지와 함께 사용해 사용자에게 유효한 쿼리 작성법을 알려주세요.

프레디킷 푸시다운 (Predicate pushdown)

WHERE 절이 있는 쿼리를 실행하면 쿼리 계획에 프레디킷 제약이 있는 ScanFilterProject 계획 노드가 포함될 수 있어요.

프레디킷 제약은 WHERE 절로 표현된 대로 단계/프래그먼트 결과에 부과된 제약의 설명이에요. 예를 들어 WHERE x > 5 AND y = 3summary 필드가 x 컬럼의 도메인은 5보다 커야 하고 y 컬럼의 도메인은 3과 같아야 한다는 제약으로 변환돼요.

쿼리 계획에 ScanFilterProject 연산이 있으면 Trino는 커넥터 메타데이터 서비스의 applyFilter 메서드를 호출해 프레디킷 제약을 커넥터로 푸시다운해 쿼리를 최적화해요. 이 메서드는 지금까지 적용된 모든 최적화가 있는 테이블 핸들을 받고, Optional.empty() 또는 기존 것에서 파생된 새 테이블 핸들을 가진 응답을 반환해요.

쿼리 옵티마이저는 최적의 쿼리 계획을 찾으면서 단일 쿼리에 대해 applyFilter를 여러 번 호출할 수 있어요. 커넥터는 이 호출에 제약을 적용할 수 없다면, 일반적으로 ScanFilterProject 푸시다운을 지원하더라도 applyFilter에서 Optional.empty()를 반환해야 해요. 제약이 이미 적용된 경우에도 커넥터는 Optional.empty()를 반환해야 해요.

제약은 다음 요소를 포함해요.

  • 컬럼과 그 도메인 사이의 매핑을 정의하는 TupleDomain. Domain은 가능한 값 목록 또는 범위 목록이며, null 허용성에 대한 정보도 담아요.
  • 함수 호출 푸시다운용 표현식.
  • 표현식의 변수에서 컬럼으로의 할당 맵.
  • (선택) 컬럼과 그 값의 맵을 테스트하는 프레디킷. applyFilter 호출이 반환된 뒤에는 보유할 수 없어요.
  • (선택) 프레디킷이 의존하는 컬럼 집합. 프레디킷이 있으면 반드시 있어야 해요.

프레디킷과 summary를 모두 사용할 수 있다면, 프레디킷은 값 필터링에서 더 엄격함이 보장되며 사용하면 쿼리 성능을 크게 높일 수 있어요.

하지만 테이블 핸들에 프레디킷을 저장해 나중에 사용하는 것은 불가능해요. applyFilter 호출이 반환된 뒤에는 프레디킷을 보유할 수 없기 때문이에요. 이는 전체 파티션을 필터링하는 데 사용되며 푸시다운되지 않아요. 대신 summary를 테이블 핸들에 저장해 푸시다운할 수 있어요.

프레디킷과 summary 사이의 이 겹침은 역사적 이유 때문이에요. 단순 비교 푸시다운이 먼저 summary를 통해 구현됐고, 더 표현력 있는 프레디킷이 필요했던 LIKE 같은 더 복잡한 필터가 나중에 추가됐기 때문이에요.

제약이 부분적으로만 푸시다운될 수 있다면(예: 범위 매칭을 지원하지 않는 데이터베이스의 커넥터를 WHERE x = 2 AND y > 5 쿼리에 사용하는 경우), y 컬럼 제약은 applyFilterConstraintApplicationResult에서 반환되어야 해요. 이 경우 y > 5 조건은 Trino에서 적용되고 푸시다운되지 않아요.

다음은 TupleDomain만 보는 간단한 예시예요.

@Override
public Optional<ConstraintApplicationResult<ConnectorTableHandle>> applyFilter(
        ConnectorSession session,
        ConnectorTableHandle tableHandle,
        Constraint constraint)
{
    ExampleTableHandle handle = (ExampleTableHandle) tableHandle;

    TupleDomain<ColumnHandle> oldDomain = handle.getConstraint();
    TupleDomain<ColumnHandle> newDomain = oldDomain.intersect(constraint.getSummary());
    if (oldDomain.equals(newDomain)) {
        // Nothing has changed, return empty Option
        return Optional.empty();
    }

    handle = new ExampleTableHandle(newDomain);
    return Optional.of(new ConstraintApplicationResult<>(handle, TupleDomain.all(), false));
}

제약의 TupleDomain은 이미 TableHandle에 적용된 TupleDomain과 교집합되어 newDomain을 형성해요. 필터링이 바뀌지 않았다면 이 최적화 경로가 끝에 도달했다는 것을 플래너에 알리기 위해 Optional.empty() 결과가 반환돼요.

이 예시에서 커넥터는 데이터 소스에서 같은 의미로 지원되는 모든 Trino 데이터 타입과 함께 TupleDomain을 푸시다운해요. 그 결과 Trino에서는 필터가 필요 없고, ConstraintApplicationResultremainingFilterTupleDomain.all()로 설정해요.

이 푸시다운 구현은 MongoMetadata, BigQueryMetadata, KafkaMetadata를 포함한 많은 Trino 커넥터와 상당히 비슷해요.

다음의 더 복잡한 예시는 기반 데이터 소스에서 직접 사용할 수 없어 매핑해야 하는 Trino 데이터 타입을 보여줘요.

@Override
public Optional<ConstraintApplicationResult<ConnectorTableHandle>> applyFilter(
        ConnectorSession session,
        ConnectorTableHandle table,
        Constraint constraint)
{
    JdbcTableHandle handle = (JdbcTableHandle) table;

    TupleDomain<ColumnHandle> oldDomain = handle.getConstraint();
    TupleDomain<ColumnHandle> newDomain = oldDomain.intersect(constraint.getSummary());
    TupleDomain<ColumnHandle> remainingFilter;
    if (newDomain.isNone()) {
        newConstraintExpressions = ImmutableList.of();
        remainingFilter = TupleDomain.all();
        remainingExpression = Optional.of(Constant.TRUE);
    }
    else {
        // We need to decide which columns to push down.
        // Since this is a base class for many JDBC-based connectors, each
        // having different Trino type mappings and comparison semantics
        // it needs to be flexible.

        Map<ColumnHandle, Domain> domains = newDomain.getDomains().orElseThrow();
        List<JdbcColumnHandle> columnHandles = domains.keySet().stream()
                .map(JdbcColumnHandle.class::cast)
                .collect(toImmutableList());

        // Get information about how to push down every column based on its
        // JDBC data type
        List<ColumnMapping> columnMappings = jdbcClient.toColumnMappings(
                session,
                columnHandles.stream()
                        .map(JdbcColumnHandle::getJdbcTypeHandle)
                        .collect(toImmutableList()));

        // Calculate the domains which can be safely pushed down (supported)
        // and those which need to be filtered in Trino (unsupported)
        Map<ColumnHandle, Domain> supported = new HashMap<>();
        Map<ColumnHandle, Domain> unsupported = new HashMap<>();
        for (int i = 0; i < columnHandles.size(); i++) {
            JdbcColumnHandle column = columnHandles.get(i);
            DomainPushdownResult pushdownResult =
                columnMappings.get(i).getPredicatePushdownController().apply(
                    session,
                    domains.get(column));
            supported.put(column, pushdownResult.getPushedDown());
            unsupported.put(column, pushdownResult.getRemainingFilter());
        }

        newDomain = TupleDomain.withColumnDomains(supported);
        remainingFilter = TupleDomain.withColumnDomains(unsupported);
    }

    // Return empty Optional if nothing changed in filtering
    if (oldDomain.equals(newDomain)) {
        return Optional.empty();
    }

    handle = new JdbcTableHandle(
            handle.getRelationHandle(),
            newDomain,
            ...);

    return Optional.of(
            new ConstraintApplicationResult<>(
                handle,
                remainingFilter));
}

이 예시는 여러 JDBC 호환 데이터 소스의 특정 요구 사항을 처리하면서 많은 JDBC 커넥터의 기반 클래스를 구현하는 것을 보여줘요. 푸시다운된 제약이 기반 데이터 소스에서 정확히 똑같이 동작하고 Trino에서와 같은 결과를 만드는 것을 보장해요. 예를 들어 문자열 비교가 대소문자를 구분하지 않는 데이터베이스에서는, Trino의 문자열 비교 연산이 대소문자를 구분하므로 푸시다운이 동작하지 않아요.

PredicatePushdownController 인터페이스는 JDBC 호환 데이터 소스에서 컬럼 도메인이 푸시다운될 수 있는지 결정해요. 위 예시에서는 그 데이터베이스에 특화된 JdbcClient 구현에서 호출돼요. JDBC 호환 데이터 소스가 아닌 곳에서는 타입 기반 푸시다운이 PredicatePushdownController 인터페이스를 거치지 않고 직접 구현돼요.

다음 예시는 세션 플래그로 활성화되는 표현식 푸시다운을 추가해요.

@Override
public Optional<ConstraintApplicationResult<ConnectorTableHandle>> applyFilter(
        ConnectorSession session,
        ConnectorTableHandle table,
        Constraint constraint)
{
    JdbcTableHandle handle = (JdbcTableHandle) table;

    TupleDomain<ColumnHandle> oldDomain = handle.getConstraint();
    TupleDomain<ColumnHandle> newDomain = oldDomain.intersect(constraint.getSummary());
    List<String> newConstraintExpressions;
    TupleDomain<ColumnHandle> remainingFilter;
    Optional<ConnectorExpression> remainingExpression;
    if (newDomain.isNone()) {
        newConstraintExpressions = ImmutableList.of();
        remainingFilter = TupleDomain.all();
        remainingExpression = Optional.of(Constant.TRUE);
    }
    else {
        // We need to decide which columns to push down.
        // Since this is a base class for many JDBC-based connectors, each
        // having different Trino type mappings and comparison semantics
        // it needs to be flexible.

        Map<ColumnHandle, Domain> domains = newDomain.getDomains().orElseThrow();
        List<JdbcColumnHandle> columnHandles = domains.keySet().stream()
                .map(JdbcColumnHandle.class::cast)
                .collect(toImmutableList());

        // Get information about how to push down every column based on its
        // JDBC data type
        List<ColumnMapping> columnMappings = jdbcClient.toColumnMappings(
                session,
                columnHandles.stream()
                        .map(JdbcColumnHandle::getJdbcTypeHandle)
                        .collect(toImmutableList()));

        // Calculate the domains which can be safely pushed down (supported)
        // and those which need to be filtered in Trino (unsupported)
        Map<ColumnHandle, Domain> supported = new HashMap<>();
        Map<ColumnHandle, Domain> unsupported = new HashMap<>();
        for (int i = 0; i < columnHandles.size(); i++) {
            JdbcColumnHandle column = columnHandles.get(i);
            DomainPushdownResult pushdownResult =
                columnMappings.get(i).getPredicatePushdownController().apply(
                    session,
                    domains.get(column));
            supported.put(column, pushdownResult.getPushedDown());
            unsupported.put(column, pushdownResult.getRemainingFilter());
        }

        newDomain = TupleDomain.withColumnDomains(supported);
        remainingFilter = TupleDomain.withColumnDomains(unsupported);

        // Do we want to handle expression pushdown?
        if (isComplexExpressionPushdown(session)) {
            List<String> newExpressions = new ArrayList<>();
            List<ConnectorExpression> remainingExpressions = new ArrayList<>();
            // Each expression can be broken down into a list of conjuncts
            // joined with AND. We handle each conjunct separately.
            for (ConnectorExpression expression : extractConjuncts(constraint.getExpression())) {
                // Try to convert the conjunct into something which is
                // understood by the underlying JDBC data source
                Optional<String> converted = jdbcClient.convertPredicate(
                    session,
                    expression,
                    constraint.getAssignments());
                if (converted.isPresent()) {
                    newExpressions.add(converted.get());
                }
                else {
                    remainingExpressions.add(expression);
                }
            }
            // Calculate which parts of the expression can be pushed down
            // and which need to be calculated in Trino engine
            newConstraintExpressions = ImmutableSet.<String>builder()
                    .addAll(handle.getConstraintExpressions())
                    .addAll(newExpressions)
                    .build().asList();
            remainingExpression = Optional.of(and(remainingExpressions));
        }
        else {
            newConstraintExpressions = ImmutableList.of();
            remainingExpression = Optional.empty();
        }
    }

    // Return empty Optional if nothing changed in filtering
    if (oldDomain.equals(newDomain) &&
            handle.getConstraintExpressions().equals(newConstraintExpressions)) {
        return Optional.empty();
    }

    handle = new JdbcTableHandle(
            handle.getRelationHandle(),
            newDomain,
            newConstraintExpressions,
            ...);

    return Optional.of(
            remainingExpression.isPresent()
                    ? new ConstraintApplicationResult<>(
                        handle,
                        remainingFilter,
                        remainingExpression.get())
                    : new ConstraintApplicationResult<>(
                        handle,
                        remainingFilter));
}

ConnectorExpressionTupleDomain과 비슷하게 분할돼요. 각 표현식은 독립적인 conjunct로 쪼갤 수 있어요. Conjunct는 AND 연산자로 연결하면 원래 표현식과 동일한 더 작은 표현식이에요. 각 conjunct는 개별적으로 처리할 수 있어요. 각각은 더 유연하게 JdbcClient 구현이 정의하는 커넥터별 규칙으로 변환돼요. 변환되지 않은 conjunct는 remainingExpression으로 반환되고 Trino 엔진이 평가해요.

ConnectorSplitManager

스플릿 매니저는 테이블의 데이터를 Trino가 워커에 분배해 처리할 개별 청크로 나눠요. 예를 들어 Hive 커넥터는 각 Hive 파티션의 파일을 나열하고 파일마다 하나 이상의 스플릿을 만들어요. 파티션된 데이터가 없는 데이터 소스는 테이블 전체에 대해 단일 스플릿을 반환하는 것이 좋은 전략이에요. 예제 HTTP 커넥터가 사용하는 전략이에요.

ConnectorRecordSetProvider

스플릿, 테이블 핸들, 컬럼 목록이 주어지면 레코드 셋 프로바이더는 Trino 실행 엔진에 데이터를 전달하는 책임을 가져요.

테이블과 컬럼 핸들은 가상 테이블을 나타내요. 이들은 커넥터의 메타데이터 서비스가 만들며, Trino가 쿼리 계획·최적화 중에 호출해요. 이런 가상 테이블은 커넥터 데이터 소스의 단일 컬렉션에 직접 매핑될 필요는 없어요. 커넥터가 푸시다운을 지원하면 다른 가상 테이블에서 파생된 여러 가상 테이블이 있어 기반 데이터의 다른 뷰를 제시할 수 있어요.

프로바이더는 RecordSet을 만들고, 이는 다시 Trino가 각 행의 컬럼 값을 읽는 데 사용하는 RecordCursor를 만들어요.

제공된 레코드 셋은 ConnectorRecordSetProvider.getRecordSet() 메서드로 전달된 컬럼 핸들 목록과 순서가 일치하는 요청된 컬럼만 포함해야 해요. 레코드 셋은 TableScan 연산과 연결된 TableHandle이 나타내는 "가상 테이블" 안의 모든 행을 반환해야 해요.

성능이 중요하지 않은 단순한 커넥터의 경우 레코드 셋 프로바이더는 InMemoryRecordSet 인스턴스를 반환할 수 있어요. 인메모리 레코드 셋은 모든 행에 대한 값 목록으로 만들 수 있어, RecordCursor를 구현하는 것보다 간단할 수 있어요.

RecordCursor 구현은 현재 레코드를 추적해야 해요. 컬럼 정의와 일치하는 데이터 타입으로, 숫자 위치로 컬럼 값을 반환해요. 엔진이 현재 레코드 읽기를 마치면 커서에서 advanceNextPosition을 호출해요.

타입 매핑 (Type mapping)

내장 SQL 데이터 타입은 캐리어 타입으로 서로 다른 Java 타입을 사용해요.

SQL 타입 Java 타입
BOOLEAN boolean
TINYINT long
SMALLINT long
INTEGER long
BIGINT long
REAL long
DOUBLE double
DECIMAL 정밀도 19 이하까지는 long; 정밀도 19 초과는 Int128
VARCHAR Slice
CHAR Slice
VARBINARY Slice
JSON Slice
DATE long
TIME(P) long
TIME WITH TIME ZONE 정밀도 9까지는 long; 정밀도 9 초과는 LongTimeWithTimeZone
TIMESTAMP(P) 정밀도 6까지는 long; 정밀도 6 초과는 LongTimestamp
TIMESTAMP(P) WITH TIME ZONE 정밀도 3까지는 long; 정밀도 3 초과는 LongTimestampWithTimeZone
INTERVAL YEAR TO MONTH long
INTERVAL DAY TO SECOND long
ARRAY Block
MAP Block
ROW Block
IPADDRESS Slice
UUID Slice
HyperLogLog Slice
P4HyperLogLog Slice
SetDigest Slice
QDigest Slice
TDigest TDigest

RecordCursor.getType(int field) 메서드는 필드의 SQL 타입을 반환하고, 필드 값은 캐리어 타입과 일치하는 다음 메서드 중 하나로 반환돼요.

  • getBoolean(int field)
  • getLong(int field)
  • getDouble(int field)
  • getSlice(int field)
  • getObject(int field)

real 타입의 값은 NaN 보존과 함께 IEEE 754 부동소수점 "단일 형식" 비트 레이아웃을 사용해 long으로 인코딩돼요. 이는 java.lang.Float.floatToRawIntBits 정적 메서드로 수행할 수 있어요.

일반 정밀도의 timestamp(p) with time zonetime(p) with time zone 타입의 값은 io.trino.spi.type.DateTimeEncoding 클래스의 pack()이나 packDateTimeWithZone() 같은 정적 메서드로 long으로 변환할 수 있어요.

UTF-8 인코딩 문자열은 Slices.utf8Slice() 정적 메서드로 Slice로 변환할 수 있어요.

참고

Slice 클래스는 io.airlift:slice 패키지에서 제공돼요.

Int128 객체는 Int128.valueOf() 메서드로 만들 수 있어요.

다음 예시는 array(varchar) 컬럼에 대한 블록을 만들어요.

private Block encodeArray(List<String> names)
{
    BlockBuilder builder = VARCHAR.createBlockBuilder(null, names.size());
    blockBuilder.buildEntry(elementBuilder -> names.forEach(name -> {
        if (name == null) {
            elementBuilder.appendNull();
        }
        else {
            VARCHAR.writeString(elementBuilder, name);
        }
    }));
    return builder.build();
}

다음 예시는 map(varchar, varchar) 컬럼에 대한 SqlMap 객체를 만들어요.

private SqlMap encodeMap(Map<String, ?> map)
{
    MapType mapType = typeManager.getType(TypeDescriptor.mapType(
                            VARCHAR.getTypeDescriptor(),
                            VARCHAR.getTypeDescriptor()));
    MapBlockBuilder values = mapType.createBlockBuilder(null, map != null ? map.size() : 0);
    if (map == null) {
        values.appendNull();
        return values.build().getObject(0, Block.class);
    }
    values.buildEntry((keyBuilder, valueBuilder) -> map.foreach((key, value) -> {
        VARCHAR.writeString(keyBuilder, key);
        if (value == null) {
            valueBuilder.appendNull();
        }
        else {
            VARCHAR.writeString(valueBuilder, value.toString());
        }
    }));
    return values.build().getObject(0, SqlMap.class);
}

ConnectorPageSourceProvider

스플릿, 테이블 핸들, 컬럼 목록이 주어지면 페이지 소스 프로바이더는 Trino 실행 엔진에 데이터를 전달하는 책임을 가져요. ConnectorPageSource를 만들고, 이는 다시 Trino가 컬럼 값을 읽는 데 사용하는 Page 객체를 만들어요.

구현하지 않으면 기본 RecordPageSourceProvider가 사용돼요. 레코드 셋 프로바이더가 주어지면 레코드 셋의 레코드에서 Page 객체를 만드는 RecordPageSource 인스턴스를 반환해요.

커넥터는 페이지를 직접 만드는 것이 가능할 때 레코드 셋 프로바이더 대신 페이지 소스 프로바이더를 구현해야 해요. 레코드 셋 프로바이더의 개별 레코드를 페이지로 변환하면 쿼리 실행 중 오버헤드가 추가돼요.

ConnectorPageSinkProvider

insert 테이블 핸들이 주어지면 페이지 싱크 프로바이더는 Trino 실행 엔진에서 데이터를 소비하는 책임을 가져요. ConnectorPageSink를 만들고, 이는 다시 컬럼 값을 담은 Page 객체를 받아요.

페이지를 반복해 개별 값에 접근하는 방법을 보여주는 예시:

@Override
public CompletableFuture<?> appendPage(Page page)
{
    for (int channel = 0; channel < page.getChannelCount(); channel++) {
        Block block = page.getBlock(channel);
        for (int position = 0; position < page.getPositionCount(); position++) {
            if (block.isNull(position)) {
                // or handle this differently
                continue;
            }

            // channel should match the column number in the table
            // use it to determine the expected column type
            String value = VARCHAR.getSlice(block, position).toStringUtf8();
            // TODO do something with the value
        }
    }
    return NOT_BLOCKED;
}

참고

커넥터를 Trino와 함께 새 의존성으로 패키징할 때 core/trino-server/src/main/provisio/trino.xml에 아티팩트를 등록하세요.

더 알아보기 (Learn more)

실제 커넥터 구현 예시는 예제 HTTP 커넥터예제 JDBC 커넥터를 참고해 보세요. 커넥터가 제공하는 서비스 전반은 SPI 개요에서 볼 수 있어요.