수집 변환

수집 변환 (Ingestion Transformations)

원본 소스 데이터는 Pinot에 들어가기 전에 여러 변환이 필요한 경우가 많아요.

출처: Ingestion Transformations

본문

변환에는 중첩 객체에서 레코드 추출, 특정 컬럼에 대한 간단한 변환 함수 적용, 원치 않는 컬럼 필터링, 데이터셋 간 조인 같은 고급 작업이 포함됩니다.

이런 작업을 수행하려면 보통 전처리 작업이 필요합니다. 스트리밍 데이터 소스에서는 Samza 작업을 작성하고 변환된 데이터를 저장할 중간 토픽을 만들 수도 있습니다.

간단한 변환의 경우 이는 배치/스트림 데이터 소스의 불일치를 초래하고 유지·운영 오버헤드를 높일 수 있습니다.

이를 쉽게 하기 위해 Pinot는 table config을 통해 적용할 수 있는 변환을 지원합니다.

{% hint style="warning" %} 실시간 컨슈머가 실행 중일 때 변환된 컬럼을 추가하거나 변경하면, 현재 소비 중인 세그먼트는 커밋할 때까지 이전 변환 계획을 계속 사용할 수 있습니다. 수집 중 새 컬럼 추가(#add-a-new-column-during-ingestion) 단계를 사용하세요 — transform 변경에는 pause를 권장하지만, 단순 기본값 전용 컬럼은 보통 reload나 forceCommit만으로 충분합니다. {% endhint %}

변환 함수

Pinot는 다음 함수를 지원합니다:

  • Groovy 함수
  • 내장 함수

{% hint style="warning" %} 하나의 변환 함수는 Groovy와 내장 함수를 섞을 수 없습니다; 한 번에 한 종류의 함수만 사용하세요. {% endhint %}

Groovy 함수

Groovy 함수는 다음 문법으로 정의할 수 있습니다:

Groovy({groovy script}, argument1, argument2...argumentN)

유효한 Groovy 표현식은 무엇이든 사용할 수 있습니다.

:warning: Groovy 활성화 {: .warning}

수집 변환에서 실행 가능한 Groovy를 허용하는 것은 보안 취약점이 될 수 있습니다. 수집용 Groovy를 활성화하려면 다음 컨트롤러 설정을 지정하세요:

controller.disable.ingestion.groovy=false

설정하지 않으면 수집 변환용 Groovy는 기본적으로 비활성화됩니다.

Groovy 정적 분석

Pinot는 Groovy 스크립트를 컴파일하기 전에 정적 분석을 적용할 수도 있습니다. 이는 클러스터 수준 Groovy 분석기 설정으로 구성됩니다:

  • pinot.groovy.all.static.analyzer: 쿼리 시간과 수집 시간 Groovy 모두에 대한 기본 분석기
  • pinot.groovy.ingestion.static.analyzer: Pinot가 테이블 설정에서 Groovy를 검증할 때 사용하는 수집 전용 오버라이드
  • pinot.groovy.query.static.analyzer: Pinot가 브로커 쿼리에서 Groovy를 검증할 때 사용하는 쿼리 전용 오버라이드

수집 전용이나 쿼리 전용 설정이 없으면 Pinot는 pinot.groovy.all.static.analyzer로 폴백합니다. 이 클러스터 설정이 모두 없으면 Groovy 정적 분석은 비활성화됩니다.

각 정적 분석기 설정은 다음 필드를 가진 JSON 객체입니다:

  • allowedReceivers: Groovy 메서드 호출이 대상으로 삼을 수 있는 정규화된 수신자 클래스
  • allowedImports: 스크립트에서 허용되는 정규화된 임포트
  • allowedStaticImports: 스크립트에서 허용되는 정규화된 정적 임포트
  • disallowedMethodNames: 수신자가 허용되더라도 Pinot가 거부하는 메서드 이름
  • methodDefinitionAllowed: 스크립트 안에서 Groovy 메서드 정의가 허용되는지 여부

정적 분석 자체는 Groovy를 활성화하지 않습니다. 수집 변환에서 Groovy를 쓰려면 여전히 controller.disable.ingestion.groovy=false가 필요합니다.

컨트롤러 API로 이러한 설정을 조회·업데이트하세요:

  • GET /cluster/configs/groovy/staticAnalyzerConfig/default는 Pinot의 내장 샘플 설정을 반환합니다.
  • GET /cluster/configs/groovy/staticAnalyzerConfig는 클러스터의 현재 Groovy 정적 분석기 설정 오버라이드를 반환합니다.
  • POST /cluster/configs/groovy/staticAnalyzerConfig는 하나 이상의 Groovy 분석기 설정을 업데이트합니다.

이 엔드포인트의 전체 요청·응답 예제는 controller API reference에 있습니다.

내장 Pinot 함수

Pinot는 수집 변환용 등록된 스칼라 함수를 지원합니다. 함수는 @ScalarFunction으로 주석 처리됩니다. 예: toEpochSeconds.

함수 변동성 요구사항

REALTIME 테이블의 경우, 새로 추가되거나 변경된 ingestionConfig.transformConfigs는 명시적 인자에만 의존하는 스칼라 함수(IMMUTABLE 함수)만 사용할 수 있습니다. OFFLINE 테이블 수준 변환은 모든 변동성 범주의 등록된 스칼라 함수를 사용할 수 있습니다. 이로 인해 배치 수집이 입력 파일에 없는 수집 타임스탬프 같은 값을 now()로 생성할 수 있습니다.

스키마 수준 transformFunction 정의와 upsertConfig.postPartialUpsertTransformConfigs는 테이블 타입과 무관하게 기존 변동성 검사를 유지합니다.

불변 함수가 필요한 곳에서는 중첩 스칼라 함수 호출도 검사에 포함됩니다. 예를 들어 now(), ago(), agoMV(), 시드 없는 rand()는 같은 입력에 대해 다른 값을 만들 수 있으므로 거부됩니다. rand(seed)는 명시적 시드에서 결과가 유도되므로 유효합니다.

Pinot는 호환성을 위해 기존 변환 정의를 그대로 실행합니다. 불변성이 필요한 변환 표면에서는 스키마나 테이블 설정을 제출하기 전에 불변 함수를 입력 유도 대안으로 교체하세요. 이 검증은 함수의 쿼리 시간 동작을 바꾸지 않습니다. Groovy 표현식은 스칼라 함수 변동성 메타데이터가 아니라 기존 Groovy 검증을 계속 통과합니다.

아래는 수집 변환에 흔히 쓰이는 내장 Pinot 함수 일부입니다.

DateTime 함수

이 함수들은 시간 변환을 가능하게 합니다.

toEpochXXX

epoch 밀리초를 더 높은 정밀도로 변환합니다.

함수 이름 설명
toEpochSeconds epoch 밀리초를 epoch 초로 변환. 사용법: "toEpochSeconds(millis)"
toEpochMinutes epoch 밀리초를 epoch 분으로 변환. 사용법: "toEpochMinutes(millis)"
toEpochHours epoch 밀리초를 epoch 시간으로 변환. 사용법: "toEpochHours(millis)"
toEpochDays epoch 밀리초를 epoch 일로 변환. 사용법: "toEpochDays(millis)"

toEpochXXXRounded

epoch 밀리초를 다른 정밀도로 변환하되 가장 가까운 반올림 버킷으로 반올림합니다. 예를 들어 1588469352000은 26474489 minutesSinceEpoch입니다. toEpochMinutesRounded(1588469352000) = 26474480

함수 이름 설명
toEpochSecondsRounded epoch 밀리초를 epoch 초로 변환, 가장 가까운 반올림 버킷으로 반올림 "toEpochSecondsRounded(millis, 30)"
toEpochMinutesRounded epoch 밀리초를 epoch 분으로 변환, 가장 가까운 반올림 버킷으로 반올림 "toEpochMinutesRounded(millis, 10)"
toEpochHoursRounded epoch 밀리초를 epoch 시간으로 변환, 가장 가까운 반올림 버킷으로 반올림 "toEpochHoursRounded(millis, 6)"
toEpochDaysRounded epoch 밀리초를 epoch 일로 변환, 가장 가까운 반올림 버킷으로 반올림 "toEpochDaysRounded(millis, 7)"

fromEpochXXX

epoch 정밀도를 밀리초로 변환합니다.

함수 이름 설명
fromEpochSeconds epoch 초에서 밀리초로 변환 "fromEpochSeconds(secondsSinceEpoch)"
fromEpochMinutes epoch 분에서 밀리초로 변환 "fromEpochMinutes(minutesSinceEpoch)"
fromEpochHours epoch 시간에서 밀리초로 변환 "fromEpochHours(hoursSinceEpoch)"
fromEpochDays epoch 일에서 밀리초로 변환 "fromEpochDays(daysSinceEpoch)"

간단한 날짜 형식

제공된 패턴 문자열에 따라 간단한 날짜 형식 문자열을 밀리초로(및 그 반대로) 변환합니다.

함수 이름 설명
ToDateTime 제공된 패턴에 따라 밀리초를 형식화된 날짜-시간 문자열로 변환 "toDateTime(millis, 'yyyy-MM-dd')"
FromDateTime 제공된 패턴에 따라 형식화된 날짜-시간 문자열을 밀리초로 변환 "fromDateTime(dateTimeStr, 'EEE MMM dd HH:mm:ss ZZZ yyyy')"

{% hint style="info" %} 참고

Simple Date Time 전설에 속하지 않는 문자는 이스케이프해야 합니다. 예:

"transformFunction": "fromDateTime(dateTimeStr, 'yyyy-MM-dd''T''HH:mm:ss')" {% endhint %}

JSON 함수

함수 이름 설명
json_format JSON/AVRO 복합 객체를 문자열로 변환합니다. 이 JSON 맵은 이후 jsonExtractScalar 함수로 조회할 수 있습니다. "json_format(jsonMapField)"

지리공간 함수

지리공간(scalar) 함수를 수집 변환으로 사용할 수 있습니다. 수집 시점에 점을 구체화하고, WKT/WKB/GeoJSON을 파싱하고, geometry/geography를 변환하고, H3 그리드 id를 계산하는 데 사용하세요. Geometry와 geography 값은 스키마에서 BYTES로 저장됩니다.

함수 이름 반환 설명
ST_Point / stPoint BYTES (Point) x/y 좌표에서 점을 만듭니다. 선택적 세 번째 인자는 geography 여부를 나타냅니다.
ST_Polygon BYTES (Polygon) 폴리곤 WKT를 평면 폴리곤 geometry로 파싱합니다.
ST_GeomFromText / stGeomFromText BYTES (Geometry) WKT 문자열을 geometry로 파싱합니다.
ST_GeogFromText / stGeogFromText BYTES (Geography) WKT 문자열을 geography로 파싱합니다.
ST_GeomFromWKB / stGeomFromWKB BYTES (Geometry) WKB 바이트를 geometry로 파싱합니다.
ST_GeogFromWKB / stGeogFromWKB BYTES (Geography) WKB 바이트를 geography로 파싱합니다.
ST_GeomFromGeoJSON / stGeomFromGeoJson BYTES (Geometry) GeoJSON 문자열을 geometry로 파싱합니다.
ST_GeogFromGeoJSON / stGeogFromGeoJson BYTES (Geography) GeoJSON 문자열을 geography로 파싱합니다.
ST_AsText / stAsText STRING geometry/geography를 WKT로 직렬화합니다.
ST_AsBinary / stAsBinary BYTES geometry/geography를 WKB로 직렬화합니다.
ST_AsGeoJSON / stAsGeoJson STRING geometry/geography를 GeoJSON으로 직렬화합니다.
ST_GeometryType / stGeometryType STRING geometry 타입 이름(Point, Polygon 등)을 반환합니다.
ST_Area / stArea DOUBLE geometry의 평면 면적 또는 geography의 구면 면적(제곱미터)을 계산합니다.
ST_Distance / stDistance DOUBLE 두 geometry(직교) 또는 geography(미터) 사이의 거리.
ST_Contains / stContains INT (0/1) 첫 번째 geometry/geography가 두 번째를 포함하는지 여부.
ST_Equals / stEquals INT (0/1) 두 geometry가 같은지 여부.
ST_Within / stWithin INT (0/1) 첫 번째 geometry가 두 번째 안에 완전히 들어있는지 여부.
toSphericalGeography BYTES (Geography) geometry 객체를 구면 geography로 변환.
toGeometry BYTES (Geometry) 구면 geography 객체를 geometry로 변환.
geoToH3 LONG (longitude, latitude, resolution) 또는 점 BYTES + resolution에서 H3 셀 id.
gridDistance LONG 두 H3 인덱스 사이의 H3 그리드 거리.
gridDisk LONG[] 원점 인덱스에서 k 그리드 스텝 내의 H3 인덱스.

{% hint style="info" %} 변환에서 함수 이름은 대소문자를 구분하지 않습니다. 전체 지리공간 함수 레퍼런스(쿼리 시간 사용법과 ST_Union 같은 집계 포함)는 Geospatial functions에 있습니다. {% endhint %}

예제: 경도/위도 컬럼에서 점과 H3 그리드 id 구체화

"ingestionConfig": {
  "transformConfigs": [
    {
      "columnName": "location",
      "transformFunction": "ST_Point(lon, lat, 1)"
    },
    {
      "columnName": "h3Index",
      "transformFunction": "geoToH3(lon, lat, 7)"
    }
  ]
}

스키마에 location을 BYTES 차원으로, h3Index를 LONG 차원으로 추가하세요. ST_Point의 세 번째 인자는 값을 평면 geometry가 아닌 geography(1/true)로 표시합니다.

예제: WKT 파싱하고 기존 점 컬럼에서 H3 유도

"ingestionConfig": {
  "transformConfigs": [
    {
      "columnName": "geom",
      "transformFunction": "ST_GeomFromText(wkt)"
    },
    {
      "columnName": "h3FromPoint",
      "transformFunction": "geoToH3(location, 9)"
    }
  ]
}

변환 유형

필터링

레코드는 수집될 때 필터링될 수 있습니다. 필터 함수는 테이블 설정의 ingestionConfigs 안 filterConfigs에 지정할 수 있습니다.

"tableConfig": {
    "tableName": ...,
    "tableType": ...,
    "ingestionConfig": {
        "filterConfig": {
            "filterFunction": "<expression>"
        }
    }
}

표현식이 true로 평가되면 레코드는 필터링됩니다. 표현식은 앞 절에서 설명한 변환 함수를 모두 사용할 수 있습니다.

timestamp 컬럼이 있는 테이블을 생각해 보세요. 타임스탬프 1589007600000보다 오래된 레코드를 걸러내려면 다음 함수를 적용할 수 있습니다:

"ingestionConfig": {
    "filterConfig": {
        "filterFunction": "Groovy({timestamp < 1589007600000}, timestamp)"
    }
}

문자열 컬럼 campaign과 다중 값 double 컬럼 prices가 있는 테이블을 생각해 보세요. campaign = 'X' 또는 'Y'이고 prices 모든 요소의 합이 100 미만인 레코드를 걸러내려면 다음 함수를 적용할 수 있습니다:

"ingestionConfig": {
    "filterConfig": {
        "filterFunction": "Groovy({(campaign == \"X\" || campaign == \"Y\") && prices.sum() < 100}, prices, campaign)"
    }
}

필터 설정은 레코드 필터링을 위한 내장 스칼라 함수의 SQL 형식 표현식도 지원합니다(v 0.11.0+부터). 예:

"ingestionConfig": {
    "filterConfig": {
        "filterFunction": "strcmp(campaign, 'X') = 0 OR strcmp(campaign, 'Y') = 0 OR timestamp < 1589007600000"
    }
}

컬럼 변환

트랜스폼 함수는 테이블 설정의 수집 설정에서 컬럼에 정의할 수 있습니다.

{% hint style="info" %} Pinot는 수집 변환을 정규화된 추출기 값에 대해 평가하지, 원시 소스 형식 인코딩에 대해서는 평가하지 않습니다. Avro, Parquet, ORC, Thrift, Protocol Buffers 같은 타입 입력 형식의 경우 부울은 Boolean로 유지되고, Byte·Short 값은 Integer로 확장되며, 논리적 날짜·시간은 java.time.LocalDate, java.time.LocalTime, java.sql.Timestamp로 나타나고, 다중 값 필드는 Object[]로, 맵·중첩 레코드는 Map<Object, Object>로 나타납니다. Pinot는 변환 단계 후 그 중간 값을 컬럼의 선언된 스키마 타입으로 강제 변환합니다.

문자열화된 부울이나 원시 epoch 날짜·시간 값 같은 형식별 특성에 의존했던 기존 변환이 있다면, 정규화된 값을 처리하거나 명시적으로 캐스팅하도록 변환을 업데이트하세요.

JSONPATHSTRING, JSONPATHSTRINGFAST, JSONPATHSTRINGFIRSTMATCH의 경우, 이미 구체화된 트리의 타입이 있는 리프가 UUID, LocalDate, LocalTime 값을 JSON 따옴표 문자열이 아닌 따옴표 없는 정규 또는 ISO-8601 문자열로 렌더링한다는 뜻입니다. 이 변경은 수집 시간 타입 객체에만 해당하며 일반 JSON 문자열 입력의 동작은 바꾸지 않습니다. {% endhint %}

{ "tableConfig": {
    "tableName": ...,
    "tableType": ...,
    "ingestionConfig": {
        "transformConfigs": [{
          "columnName": "fieldName",
          "transformFunction": "<expression>"
        }]
    },
    ...
}

예를 들어 소스 데이터에 prices와 timestamp 필드가 있다고 가정해 보겠습니다. 최대 가격을 추출해 maxPrices 필드에 저장하고, 타임스탬프를 epoch 이후 시간 수로 변환해 hoursSinceEpoch 필드에 저장하려고 합니다. 다음 변환을 적용하면 됩니다:

{
"tableName": "myTable",
...
"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "maxPrice",
      "transformFunction": "Groovy({prices.max()}, prices)" // groovy function
    },
    {
      "columnName": "hoursSinceEpoch",
      "transformFunction": "toEpochHours(timestamp)" // built-in function
    }]
  }
}

아래는 흔히 쓰는 함수 예제입니다.

문자열 연결

firstName와 lastName을 연결해 fullName을 얻습니다.

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "fullName",
      "transformFunction": "Groovy({firstName+' '+lastName}, firstName, lastName)"
    }]
}

배열에서 요소 찾기

배열 bids에서 최대값을 찾습니다.

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "maxBid",
      "transformFunction": "Groovy({bids.max{ it.toBigDecimal() }}, bids)"
    }]
}

시간 변환

timestamp를 MILLISECONDS에서 HOURS로 변환합니다.

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "hoursSinceEpoch",
      "transformFunction": "Groovy({timestamp/(1000*60*60)}, timestamp)"
    }]
}

컬럼 이름 변경

컬럼 이름을 user_id에서 userId로 바꿉니다.

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "userId",
      "transformFunction": "user_id"
    }]
}

Kafka JSON 메시지에서 필드 이름 변경

Kafka JSON 페이로드는 좋지 않은 Pinot 컬럼 이름이 되기 쉬운 키를 자주 씁니다. 흔한 예는 event-id처럼 -를 포함하는 키입니다.

transformConfigs를 사용해 소스 키를 스키마 친화적인 컬럼에 매핑하세요. 소스 키는 따옴표로 감싼 식별자로 참조합니다.

"ingestionConfig": {
  "transformConfigs": [
    {
      "columnName": "event_id",
      "transformFunction": "\"event-id\""
    },
    {
      "columnName": "event_timestamp",
      "transformFunction": "\"event-timestamp\""
    },
    {
      "columnName": "user_id",
      "transformFunction": "\"user-id\""
    }
  ]
}

{% hint style="info" %} 목적지 컬럼(예: event_id)을 Pinot 스키마에 추가하세요. {% endhint %}

공백 포함 컬럼에서 값 추출

Pinot는 공백이 있는 컬럼을 지원하지 않으므로, 소스 데이터 컬럼에 공백이 있으면 그 값을 지원되는 이름의 컬럼에 저장해야 합니다. first Name의 값을 firstName 컬럼으로 추출하려면 다음을 실행하세요:

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "firstName",
      "transformFunction": "\"first Name \""
    }]
}

삼항 연산

eventType이 IMPRESSION이면 impression을 1로 설정합니다. CLICK도 유사합니다.

"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "impressions",
      "transformFunction": "Groovy({eventType == 'IMPRESSION' ? 1: 0}, eventType)"
    },
    {
      "columnName": "clicks",
      "transformFunction": "Groovy({eventType == 'CLICK' ? 1: 0}, eventType)"
    }]
}

AVRO Map

AVRO Map을 두 개의 다중 값 컬럼으로 Pinot에 저장합니다. 매핑을 유지하려면 키를 정렬하세요.

  1. 맵의 키를 map_keys로
  2. 맵의 값을 map_values로
"ingestionConfig": {
    "transformConfigs": [{
      "columnName": "map2_keys",
      "transformFunction": "Groovy({map2.sort()*.key}, map2)"
    },
    {
      "columnName": "map2_values",
      "transformFunction": "Groovy({map2.sort()*.value}, map2)"
    }]
}

변환 연결

변환은 연결될 수 있습니다. 즉 어떤 변환이 만든 필드를 다른 변환 함수에서 사용할 수 있습니다.

예를 들어 소스 데이터의 data 필드에 다음과 같은 JSON 문서가 있다고 가정해 보겠습니다:

{
  "userId": "12345678__foo__othertext"
}

userId를 추출하는 변환 하나와, 그다음에 식별자의 숫자 부분을 꺼내는 변환을 적용할 수 있습니다:

"ingestionConfig": {
    "transformConfigs": [
      {
        "columnName": "userOid",
        "transformFunction": "jsonPathString(data, '$.userId')"
      },
      {
        "columnName": "userId",
        "transformFunction": "Groovy({Long.valueOf(userOid.substring(0, 8))}, userOid)"
      }
   ]
}

Pinot는 이러한 체인을 목록 순서가 아니라 의존성으로 검증합니다. 의존성 그래프가 유효하기만 하면, 소비자가 transformConfigs에서 더 앞에 나타나더라도 다른 변환이 만든 컬럼을 소비할 수 있습니다.

{% hint style="info" %} 중간 변환 출력은 이후 변환에서만 소비될 때는 스키마 컬럼이 필요하지 않습니다. Pinot가 저장하거나 집계하는 최종 출력은 여전히 스키마 컬럼이 필요합니다.

예를 들어 message_obj가 스키마에 없어도 다음 패턴은 유효합니다:

"ingestionConfig": {
  "transformConfigs": [
    {
      "columnName": "message_obj",
      "transformFunction": "JSONEXTRACTOBJECT(message)"
    },
    {
      "columnName": "level",
      "transformFunction": "JSONPATHSTRING(message_obj, '$.level', null)"
    },
    {
      "columnName": "msg_time",
      "transformFunction": "JSONPATHSTRING(message_obj, '$.time', null)"
    }
  ]
}

이 예에서 message_obj는 중간 컬럼이므로 Pinot가 변환 중에 그것을 구체화한 뒤 인덱싱 전에 버립니다. level과 msg_time 같은 리프 컬럼이 저장되는 출력이므로 그 컬럼들은 스키마에 있어야 합니다. {% endhint %}

평탄화

2가지 종류의 평탄화가 있습니다:

하나의 레코드를 여러 개로

이것은 아직 네이티브로 지원되지 않습니다. 이를 사용하려면 커스텀 Decoder/RecordReader를 작성할 수 있습니다. Decoder가 제공된 입력 레코드에서 여러 GenericRow를 생성하면, $MULTIPLE_RECORDS_KEY$ 키와 함께 List<GenericRow>를 목적지 GenericRow에 설정해야 합니다. 세그먼트 생성 드라이버는 이를 특수 사례로 취급해 다중 레코드 경우를 처리합니다.

복합 객체에서 속성 추출

멀티 스테이지 쿼리 엔진(MSE)은 CROSS JOIN UNNEST(...)를 사용해 쿼리 시간에 배열 컬럼을 평탄화하는 것을 지원합니다. 이는 수집 시 데이터를 미리 평탄화하지 않고 배열 필드를 확장하는 데 권장되는 접근법입니다.

SET useMultistageEngine=true;

SELECT t.id, elem
FROM myTable AS t
CROSS JOIN UNNEST(t.arrayColumn) AS u(elem)

WITH ORDINALITY, 필터링, unnested 후 집계 등 전체 문법은 Unnest operator 페이지를 참고하세요.

수집 중 새 컬럼 추가

이 절은 새 컬럼이나 변경된 컬럼이 수집 변환을 통해 채워지거나(또는 변환 로직을 바꾸는 경우) 실시간 테이블에서 채워질 때 적용됩니다. 활성 소비 세그먼트는 가변 세그먼트가 시작될 때 변환 파이프라인을 만듭니다. 세그먼트 중간에 스키마/설정을 업데이트해도 그 컨슈머에 이미 인덱싱된 행을 다시 쓰지 않으므로, 새 컨슈머를 시작할 때까지 현재 소비 중인 세그먼트의 값이 잘못되거나 없을 수 있습니다.

모든 스키마 추가에 항상 전체 pause가 필요한 것은 아닙니다. schema evolution decision table을 사용하세요:

변경 일반적인 조치
일반 컬럼 / defaultNullValue만 스키마 업데이트 → (허용되면 소비 reload가 force commit을 요청하는) reload 또는 forceCommit 사용·폴링
새 transform 또는 변경 pause → 스키마 + 테이블 설정 적용 → 완료 세그먼트 reload → resume 선호; 또는 스키마/config 적용 → forceCommit·폴링 → 이전 계획으로 커밋된 세그먼트 reload
새 컬럼의 부분-upsert 전략 스키마와 전략 설정을 함께 업데이트 → 통제된 서버 재시작 사용; 불변 핵심 upsert 설정은 새 테이블과 재수집 필요. schema evolution 참고

Transform-safe 절차 (pause 경계)

깨끗한 컨슈머 경계로 정확한 변환 값을 보장하려면:

  1. 소비를 pause합니다(pause 상태 성공 대기):

    curl -X POST {controllerHost}/tables/{tableName}/pauseConsumption
    curl -X GET  {controllerHost}/tables/{tableName}/pauseStatus
    
  2. 새 테이블 또는 스키마 설정을 적용합니다.

  3. Pinot Controller API 또는 Pinot Admin Console로 세그먼트를 reload합니다.

  4. 소비를 재개합니다:

    curl -X POST {controllerHost}/tables/{tableName}/resumeConsumption
    

대안: 긴 pause 없는 force commit

새 스키마와 테이블 설정을 적용한 뒤 현재 컨슈머를 즉시 커밋할 수 있다면:

curl -X POST {controllerHost}/tables/{tableName}/forceCommit
curl -X GET  {controllerHost}/tables/forceCommitStatus/{jobId}

numberOfSegmentsYetToBeCommitted가 0이 될 때까지 기다린 뒤, 이전 변환 계획으로 커밋된 세그먼트를 reload해 Pinot가 새 컬럼을 다시 계산하게 하세요. 세부: Force commit API.

더 알아보기 (Learn more)