CLP로 스트림 수집

CLP로 스트림 수집 (Stream Ingestion with CLP)

수집 중에 CLP로 필드를 인코딩하는 기능을 지원하는 페이지예요.

출처: Stream Ingestion with CLP

본문

{% hint style="warning" %} 이것은 실험적 기능입니다. 설정 옵션과 사용법은 안정화될 때까지 자주 변경될 수 있습니다. {% endhint %}

Kafka를 사용해 JSON 레코드를 스트림 수집할 때, 사용자는 CLP 특정 StreamMessageDecoder를 사용해 특정 필드를 CLP로 인코딩할 수 있어요.

CLP는 구조화되지 않은 로그 메시지를 더 압축 가능하게 만들면서 검색 가능성을 유지하도록 인코딩하는 압축기예요. 메시지를 세 필드로 분해하는 방식으로 동작해요:

  • 메시지의 정적 텍스트인 로그 타입(log type);
  • 반복되는 변수 값인 사전 변수(dictionary variables); 그리고
  • 비반복 변수 값 (가능하면 특별히 인코딩하므로 인코딩 변수(encoded variables)라 부름).

검색도 마찬가지로 개별 필드에 대한 쿼리로 분해돼요.

{% hint style="info" %} CLP는 로그 메시지용으로 설계되었지만, 파일 경로 같은 다른 구조화되지 않은 텍스트도 그 인코딩의 이점을 얻을 수 있습니다. {% endhint %}

예를 들어 다음 JSON 레코드를 생각해 보세요:

{
  "timestamp": 1672531200000,
  "message": "INFO Task task_12 assigned to container: [ContainerID:container_15], operation took 0.335 seconds. 8 tasks remaining.",
  "logPath": "/mnt/data/application_123/container_15/stdout"
}

사용자가 message와 logPath 필드를 CLP로 인코딩하도록 지정하면 StreamMessageDecoder가 출력:

{
  "timestamp": 1672531200000,
  "message_logtype": "INFO Task \\x12 assigned to container: [ContainerID:\\x12], operation took \\x13 seconds. \\x11 tasks remaining.",
  "message_dictionaryVars": [
    "task_12",
    "container_15"
  ],
  "message_encodedVars": [
    1801439850948198735,
    8
  ],
  "logPath_logtype": "/mnt/data/\\x12/\\x12/stdout",
  "logPath_dictionaryVars": [
    "application_123",
    "container_15"
  ],
  "logPath_encodedVars": []
}

_logtype 접미사가 붙은 필드에서 \x11은 정수 변수의 자리표시자, \x12는 사전 변수의 자리표시자, \x13은 부동소수점 변수의 자리표시자입니다. message_encoedVars에서 부동소수점 변수 0.335는 CLP의 커스텀 인코딩으로 정수로 인코딩됩니다.

나머지 모든 필드는 org.apache.pinot.plugin.inputformat.json.JSONRecordExtractor에서와 같은 방식으로 처리돼요. 구체적으로, 테이블 스키마의 필드가 각 레코드에서 추출되고 나머지 필드는 버려져요.

설정 (Configuration)

테이블 인덱스

예시처럼 message와 logPath를 인코딩하고 싶다면 tableIndexConfig에서 다음 설정을 변경/추가해요 (관련 없는 설정은 생략):

{
  "tableIndexConfig": {
    "streamConfigs": {
      "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.clplog.CLPLogMessageDecoder",
      "stream.kafka.decoder.prop.fieldsForClpEncoding": "message,logPath",
      "stream.kafka.decoder.prop.removeProcessedFields": "true"
    },
    "varLengthDictionaryColumns": [
      "message_logtype",
      "message_dictionaryVars",
      "logPath_logtype",
      "logPath_dictionaryVars"
    ]
  }
}
  • stream.kafka.decoder.prop.fieldsForClpEncoding은 CLP로 인코딩해야 하는 필드 이름들의 쉼표로 구분된 목록.
  • stream.kafka.decoder.prop.removeProcessedFields는 선택 사항. true로 설정하면 Pinot이 파생된 CLP 컬럼(<field>_logtype, <field>_dictionaryVars, <field>_encodedVars)을 쓴 후 원래 입력 필드를 제거해요. 기본값은 false이며, 원래 필드를 파생 컬럼과 함께 유지해요.
  • logtype과 사전 변수는 길이가 크게 달라질 수 있으므로 가변 길이 사전을 사용해요.

스키마

테이블 스키마에서 CLP 인코딩 필드를 다음과 같이 구성해야 해요 (관련 없는 설정은 생략):

{
  "dimensionFieldSpecs": [
    {
      "name": "message_logtype",
      "dataType": "STRING",
      "maxLength": 2147483647
    },
    {
      "name": "message_encodedVars",
      "dataType": "LONG",
      "singleValueField": false
    },
    {
      "name": "message_dictionaryVars",
      "dataType": "STRING",
      "maxLength": 2147483647,
      "singleValueField": false
    },
    {
      "name": "logpath_logtype",
      "dataType": "STRING",
      "maxLength": 2147483647
    },
    {
      "name": "logpath_encodedVars",
      "dataType": "LONG",
      "singleValueField": false
    },
    {
      "name": "logpath_dictionaryVars",
      "dataType": "STRING",
      "maxLength": 2147483647,
      "singleValueField": false
    }
  ]
}
  • logtype과 사전 변수 컬럼에 최대 가능 길이를 사용해요.
  • 사전·인코딩 변수 컬럼은 다중 값 컬럼이에요.

CLP-인코딩 필드 검색과 디코딩

CLP-인코딩 필드를 디코딩하려면 CLPDECODE를 사용해요.

CLP-인코딩 필드를 검색하려면 CLPDECODE를 LIKE와 결합할 수 있어요. 많은 수의 행을 쿼리하면 성능이 저하될 수 있다는 점에 유의하세요.

CLP-인코딩 컬럼에 대한 효율적인 검색을 또 다른 UDF로 통합하는 작업을 진행 중이에요. 이 기능의 개발은 이 설계 문서에서 추적되고 있어요.

CLP Forward Index V2

Pinot 1.3.0부터 CLP 포워드 인덱스는 V2(CLPMutableForwardIndexV2)로 업그레이드되었으며, 이제 실시간 수집 중 CLP-인코딩 컬럼의 기본이에요. 주요 개선 사항:

카디널리티 모니터링을 통한 동적 인코딩

V2는 수집 중 사전 카디널리티를 모니터링하고 인코딩 모드를 동적으로 전환해요:

  • CLP 사전 인코딩: 로그 타입과 사전 변수 카디널리티가 document 수 대비 설정 가능한 임계값 미만으로 유지될 때 사용.
  • 원시 문자열 폴백: 카디널리티가 임계값을 초과하면(docs/cardinality 비율이 10 미만으로 떨어지면), V2는 큰 사전 유지의 메모리·I/O 오버헤드를 피하기 위해 원시 문자열 포워드 인덱스로 자동 폴백.

압축 개선

V2는 V1의 무압축 고정 비트 인코딩 대신 Zstandard 청크 압축과 함께 고정 바이트 인코딩을 사용해요. 이는 대부분의 실제 로그 데이터에서 압축 비율을 크게 개선해요.

압축 코덱 옵션

fieldConfig의 compressionCodec로 CLP-인코딩 컬럼의 압축 코덱을 선택할 수 있어요:

코덱 (Codec) 설명 (Description)
CLPV2 기본 ZStandard 압축을 사용하는 CLP V2
CLPV2_ZSTD 명시적 ZStandard 압축을 사용하는 CLP V2
CLPV2_LZ4 LZ4 압축을 사용하는 CLP V2
CLP 레거시 V1 (무압축, pass-through)

예시 필드 설정:

{
  "fieldConfigList": [
    {
      "name": "message",
      "encodingType": "RAW",
      "compressionCodec": "CLPV2_ZSTD"
    }
  ]
}

불변 CLP 포워드 인덱스

가변(실시간) 세그먼트가 불변 세그먼트로 변환될 때, V2는 재인코딩 없이 가변 사전과 인덱스 데이터를 직접 복사해 V1에 있던 직렬화/역직렬화 오버헤드를 제거해요. 결과적인 불변 포워드 인덱스는 쿼리 중 효율적인 임의 접근을 위해 메모리 매핑돼요.

더 알아보기 (Learn more)