소스 input format

소스 input format (Source input formats)

Apache Druid는 JSON, CSV, 또는 TSV 같은 구분(delimited) 형식이나 어떤 커스텀 형식의 비정규화(denormalized) 데이터도 수집할 수 있어요. 문서의 대부분 예시는 JSON 형식 데이터를 사용하지만, 다른 구분 데이터를 수집하도록 Druid를 구성하는 것도 어렵지 않습니다. 새 형식에 대한 기여는 언제나 환영합니다.

이 페이지는 Druid가 지원하는 모든 기본(default) 및 코어 확장 데이터 형식을 나열해요. 커뮤니티 확장으로 지원되는 추가 데이터 형식은 커뮤니티 확장 목록을 참고하세요.

출처: 문서

본문

데이터 형식화 (Formatting data)

다음 샘플들은 Druid에서 네이티브로 지원되는 데이터 형식을 보여줘요.

JSON

{"timestamp": "2013-08-31T01:02:33Z", "page": "Gypsy Danger", "language" : "en", "user" : "nuclear", "unpatrolled" : "true", "newPage" : "true", "robot": "false", "anonymous": "false", "namespace":"article", "continent":"North America", "country":"United States", "region":"Bay Area", "city":"San Francisco", "added": 57, "deleted": 200, "delta": -143}{"timestamp": "2013-08-31T03:32:45Z", "page": "Striker Eureka", "language" : "en", "user" : "speed", "unpatrolled" : "false", "newPage" : "true", "robot": "true", "anonymous": "false", "namespace":"wikipedia", "continent":"Australia", "country":"Australia", "region":"Cantebury", "city":"Syndey", "added": 459, "deleted": 129, "delta": 330}{"timestamp": "2013-08-31T07:11:21Z", "page": "Cherno Alpha", "language" : "ru", "user" : "masterYi", "unpatrolled" : "false", "newPage" : "true", "robot": "true", "anonymous": "false", "namespace":"article", "continent":"Asia", "country":"Russia", "region":"Oblast", "city":"Moscow", "added": 123, "deleted": 12, "delta": 111}{"timestamp": "2013-08-31T11:58:39Z", "page": "Crimson Typhoon", "language" : "zh", "user" : "triplets", "unpatrolled" : "true", "newPage" : "false", "robot": "true", "anonymous": "false", "namespace":"wikipedia", "continent":"Asia", "country":"China", "region":"Shanxi", "city":"Taiyuan", "added": 905, "deleted": 5, "delta": 900}{"timestamp": "2013-08-31T12:41:27Z", "page": "Coyote Tango", "language" : "ja", "user" : "cancer", "unpatrolled" : "true", "newPage" : "false", "robot": "true", "anonymous": "false", "namespace":"wikipedia", "continent":"Asia", "country":"Japan", "region":"Kanto", "city":"Tokyo", "added": 1, "deleted": 10, "delta": -9}

CSV

2013-08-31T01:02:33Z,"Gypsy Danger","en","nuclear","true","true","false","false","article","North America","United States","Bay Area","San Francisco",57,200,-1432013-08-31T03:32:45Z,"Striker Eureka","en","speed","false","true","true","false","wikipedia","Australia","Australia","Cantebury","Syndey",459,129,3302013-08-31T07:11:21Z,"Cherno Alpha","ru","masterYi","false","true","true","false","article","Asia","Russia","Oblast","Moscow",123,12,1112013-08-31T11:58:39Z,"Crimson Typhoon","zh","triplets","true","false","true","false","wikipedia","Asia","China","Shanxi","Taiyuan",905,5,9002013-08-31T12:41:27Z,"Coyote Tango","ja","cancer","true","false","true","false","wikipedia","Asia","Japan","Kanto","Tokyo",1,10,-9

TSV (Delimited)

2013-08-31T01:02:33Z  "Gypsy Danger"  "en"  "nuclear" "true"  "true"  "false" "false" "article" "North America" "United States" "Bay Area"  "San Francisco" 57  200 -1432013-08-31T03:32:45Z  "Striker Eureka"  "en"  "speed" "false" "true"  "true"  "false" "wikipedia" "Australia" "Australia" "Cantebury" "Syndey"  459 129 3302013-08-31T07:11:21Z  "Cherno Alpha"  "ru"  "masterYi"  "false" "true"  "true"  "false" "article" "Asia"  "Russia"  "Oblast"  "Moscow"  123 12  1112013-08-31T11:58:39Z  "Crimson Typhoon" "zh"  "triplets"  "true"  "false" "true"  "false" "wikipedia" "Asia"  "China" "Shanxi"  "Taiyuan" 905 5 9002013-08-31T12:41:27Z  "Coyote Tango"  "ja"  "cancer"  "true"  "false" "true"  "false" "wikipedia" "Asia"  "Japan" "Kanto" "Tokyo" 1 10  -9

CSV와 TSV 데이터에는 컬럼 헤더가 없다는 점에 유의하세요. 데이터를 수집할 때 이를 지정할 때 이 점이 중요해집니다.

텍스트 형식 외에도 Druid는 Orc, Parquet 형식 같은 바이너리 형식도 지원해요.

커스텀 형식 (Custom formats)

Druid는 커스텀 텍스트 데이터 형식을 지원하며, Regex input format으로 파싱할 수 있어요. 다만 이렇게 데이터를 파싱하는 것은 네이티브 Java InputFormat 확장을 작성하거나 외부 스트림 프로세서를 사용하는 것보다 덜 효율적이라는 점을 알아 두세요. 새 input format의 기여를 환영합니다.

Input format

inputFormat 필드로 입력 데이터의 데이터 형식을 지정할 수 있어요.

모든 종류의 Druid 수집에는 어떤 형태의 스키마 객체가 필요합니다. 수집할 데이터의 형식은 ioConfig의 inputFormat 항목으로 지정됩니다.

JSON

JSON 데이터를 로드하도록 JSON inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 json으로 설정 | yes | | flattenSpec | JSON Object | 중첩 JSON 데이터에 대한 평탄화(flattening) 구성 지정. 자세한 내용은 flattenSpec 참고 | no | | featureSpec | JSON Object | Java용 JSON 프로세서인 Jackson이 지원하는 JSON parser features. 이 기능들은 입력 JSON 데이터의 파싱을 제어. 기능을 활성화하려면 기능 이름을 "true" Boolean 값에 매핑. 예: "featureSpec": {"ALLOW_SINGLE_QUOTES": true, "ALLOW_UNQUOTED_FIELD_NAMES": true} | no |

다음 속성들은 스트리밍 수집에서 JSON inputFormat을 사용할 때만 적용되는 특수 속성으로, 파싱 예외가 어떻게 처리되는지와 관련돼 있어요. 스트리밍 수집에서는 여러 줄 JSON 이벤트(즉 단일 JSON 이벤트가 여러 줄에 걸친 것)를 수집할 수 있습니다. 다만 파싱 예외가 발생하면 같은 스트리밍 레코드에 있는 모든 JSON 이벤트가 버려집니다.

| Field | Type | Description | Required | | assumeNewlineDelimited | Boolean | 입력이 newline 구분 JSON(각 개별 JSON 이벤트가 단일 줄에 있고 newline으로 구분)임이 알려져 있다면, 이 옵션을 true로 설정하면 더 유연한 파싱 예외 처리가 가능. 잘못된 JSON 문법의 줄만 버려지고, 유효한 JSON 이벤트가 있는 줄은 계속 수집됨 | no (기본 false) | | useJsonNodeReader | Boolean | 여러 줄 JSON 이벤트를 수집할 때 이 옵션을 활성화하면, 파싱 예외가 발생하기 전에 스트리밍 레코드 안에서 만난 유효한 JSON 이벤트를 유지하는 JSON 파서를 사용하게 됨 | no (기본 false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "json"  },  ...}

CSV

CSV 데이터를 로드하도록 CSV inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 csv로 설정 | yes | | listDelimiter | String | multi-value dimension용 커스텀 구분자 | no (기본 = ctrl+A) | | columns | JSON array | 데이터의 컬럼 지정. 데이터의 컬럼과 같은 순서여야 함 | findColumnsFromHeader가 false이거나 없으면 yes | | findColumnsFromHeader | Boolean | 설정하면 태스크가 헤더 행에서 컬럼 이름을 찾음. skipHeaderRows가 헤더에서 컬럼 이름을 찾기 전에 적용된다는 점에 유의. 예를 들어 skipHeaderRows를 2로, findColumnsFromHeader를 true로 설정하면 태스크는 처음 두 줄을 건너뛰고 세 번째 줄에서 컬럼 정보를 추출. true로 설정하면 columns는 무시됨 | no (기본 = columns가 설정되면 false, 그 외 null) | | skipHeaderRows | Integer | 설정하면 태스크가 처음 skipHeaderRows 행을 건너뜀 | no (기본 = 0) | | tryParseNumbers | Boolean | 설정하면 태스크가 숫자 문자열을 long 또는 double 데이터 타입으로(이 순서대로) 파싱하려 시도. 이 파싱은 listDelimiter로 구분된 값에도 적용됨. 값을 숫자로 파싱할 수 없으면 문자열로 유지 | no (기본 = false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "csv",    "columns" : ["timestamp","page","language","user","unpatrolled","newPage","robot","anonymous","namespace","continent","country","region","city","added","deleted","delta"]  },  ...}

TSV (Delimited)

TSV 데이터를 로드하도록 TSV inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 tsv로 설정 | yes | | delimiter | String | 데이터 값용 커스텀 구분자 | no (기본 = \t) | | listDelimiter | String | multi-value dimension용 커스텀 구분자 | no (기본 = ctrl+A) | | columns | JSON array | 데이터의 컬럼 지정. 데이터의 컬럼과 같은 순서여야 함 | findColumnsFromHeader가 false이거나 없으면 yes | | findColumnsFromHeader | Boolean | 설정하면 태스크가 헤더 행에서 컬럼 이름을 찾음. skipHeaderRows가 헤더에서 컬럼 이름을 찾기 전에 적용된다는 점에 유의. 예를 들어 skipHeaderRows를 2로, findColumnsFromHeader를 true로 설정하면 태스크는 처음 두 줄을 건너뛰고 세 번째 줄에서 컬럼 정보를 추출. true로 설정하면 columns는 무시됨 | no (기본 = columns가 설정되면 false, 그 외 null) | | skipHeaderRows | Integer | 설정하면 태스크가 처음 skipHeaderRows 행을 건너뜀 | no (기본 = 0) | | tryParseNumbers | Boolean | 설정하면 태스크가 숫자 문자열을 long 또는 double 데이터 타입으로(이 순서대로) 파싱하려 시도. 이 파싱은 listDelimiter로 구분된 값에도 적용됨. 값을 숫자로 파싱할 수 없으면 문자열로 유지 | no (기본 = false) |

delimiter를 데이터에 맞는 구분자로 반드시 바꾸세요. CSV처럼, 컬럼과 인덱싱할 컬럼의 부분 집합을 지정해야 합니다.

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "tsv",    "columns" : ["timestamp","page","language","user","unpatrolled","newPage","robot","anonymous","namespace","continent","country","region","city","added","deleted","delta"],    "delimiter":"|"  },  ...}

Lines

각 줄이 단일 필드로 취급되는 줄 지향(line-oriented) 데이터를 로드하도록 Lines inputFormat을 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 lines로 설정 | yes |

Lines input format은 입력의 각 줄을 UTF-8 텍스트로 읽고, 전체 줄을 문자열로 담은 line이라는 단일 컬럼을 만듭니다. 이는 나중에 처리하기 위해 간단한 형태로 줄 지향 데이터를 읽는 데 유용해요.

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "lines"  },  ...}

ORC

ORC input format을 사용하려면 Druid Orc 확장(druid-orc-extensions)을 로드하세요.

ORC 데이터를 로드하도록 ORC inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 orc로 설정 | yes | | flattenSpec | JSON Object | 중첩 ORC 데이터에 대한 평탄화 구성 지정. 'path' 표현식만 지원됨('jq'와 'tree'는 사용 불가). 자세한 내용은 flattenSpec 참고 | no | | binaryAsString | Boolean | 논리적으로 문자열로 표시되지 않은 binary orc 컬럼을 UTF-8 인코딩 문자열로 취급할지 지정 | no (기본 = false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "orc",    "flattenSpec": {      "useFieldDiscovery": true,      "fields": [        {          "type": "path",          "name": "nested",          "expr": "$.path.to.nested"        }      ]    },    "binaryAsString": false  },  ...}

Parquet

Parquet input format을 사용하려면 Druid Parquet 확장(druid-parquet-extensions)을 로드하세요.

Parquet 데이터를 로드하도록 Parquet inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 parquet로 설정 | yes | | flattenSpec | JSON Object | Parquet 파일에서 중첩 값을 추출할 flattenSpec 정의. 'path' 표현식만 지원됨('jq'와 'tree'는 사용 불가) | no (기본은 'root' 레벨 속성 자동 발견) | | binaryAsString | Boolean | 논리적으로 문자열 또는 enum 타입으로 표시되지 않은 bytes parquet 컬럼을 UTF-8 인코딩 문자열로 취급할지 지정 | no (기본 = false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "parquet",    "flattenSpec": {      "useFieldDiscovery": true,      "fields": [        {          "type": "path",          "name": "nested",          "expr": "$.path.to.nested"        }      ]    },    "binaryAsString": false  },  ...}

Avro Stream

Avro Stream input format을 사용하려면 Druid Avro 확장(druid-avro-extensions)을 로드하세요.

Druid가 Avro 타입을 다루는 방법에 대한 자세한 내용은 Avro Types 섹션을 참고하세요.

Avro 데이터를 로드하도록 Avro inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 avro_stream으로 설정 | yes | | flattenSpec | JSON Object | Avro 레코드에서 중첩 값을 추출할 flattenSpec 정의. 'path' 표현식만 지원됨('jq'는 사용 불가) | no (기본은 'root' 레벨 속성 자동 발견) | | avroBytesDecoder | JSON Object | 바이트를 Avro 레코드로 디코딩하는 방법 지정 | yes | | binaryAsString | Boolean | 논리적으로 문자열 또는 enum 타입으로 표시되지 않은 bytes Avro 컬럼을 UTF-8 인코딩 문자열로 취급할지 지정 | no (기본 = false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "avro_stream",    "avroBytesDecoder": {      "type": "schema_inline",      "schema": {        //your schema goes here, for example        "namespace": "org.apache.druid.data",        "name": "User",        "type": "record",        "fields": [          { "name": "FullName", "type": "string" },          { "name": "Country", "type": "string" }        ]      }    },    "flattenSpec": {      "useFieldDiscovery": true,      "fields": [        {          "type": "path",          "name": "someRecord_subInt",          "expr": "$.someRecord.subInt"        }      ]    },    "binaryAsString": false  },  ...}
Avro Bytes Decoder

type이 포함되지 않으면 avroBytesDecoder는 schema_repo로 기본 설정됩니다.

Inline Schema 기반 Avro Bytes Decoder

정보

"schema_inline" 디코더는 고정 스키마로 Avro 레코드를 읽으며 스키마 마이그레이션을 지원하지 않아요. 나중에 스키마를 마이그레이션해야 할 수도 있다면, 파서가 레코드를 읽을 올바른 Avro 스키마를 식별하게 해 주는 메시지 헤더를 사용하는 다른 디코더 중 하나를 고려하세요.

이 디코더는 모든 입력 이벤트를 같은 스키마로 읽을 수 있을 때 사용할 수 있어요. 이 경우 아래에서 설명하듯 input task JSON 자체에 스키마를 지정합니다.

..."avroBytesDecoder": {  "type": "schema_inline",  "schema": {    //your schema goes here, for example    "namespace": "org.apache.druid.data",    "name": "User",    "type": "record",    "fields": [      { "name": "FullName", "type": "string" },      { "name": "Country", "type": "string" }    ]  }}...
Multiple Inline Schemas 기반 Avro Bytes Decoder

서로 다른 입력 이벤트가 서로 다른 읽기 스키마를 가질 수 있을 때 이 디코더를 사용하세요. 이 경우 아래에서 설명하듯 input task JSON 자체에 스키마를 지정합니다.

..."avroBytesDecoder": {  "type": "multiple_schemas_inline",  "schemas": {    //your id -> schema map goes here, for example    "1": {      "namespace": "org.apache.druid.data",      "name": "User",      "type": "record",      "fields": [        { "name": "FullName", "type": "string" },        { "name": "Country", "type": "string" }      ]    },    "2": {      "namespace": "org.apache.druid.otherdata",      "name": "UserIdentity",      "type": "record",      "fields": [        { "name": "Name", "type": "string" },        { "name": "Location", "type": "string" }      ]    },    ...    ...  }}...

이것은 본질적으로 정수 스키마 ID를 Avro 스키마 객체로 매핑한 맵이에요. 이 파서는 레코드가 다음 형식을 가진다고 가정합니다. 처음 1바이트는 버전이며 항상 1이어야 합니다. 다음 4바이트는 big-endian 바이트 순서로 직렬화된 정수 스키마 ID입니다. 나머지 바이트는 직렬화된 avro 메시지를 담고 있습니다.

SchemaRepo 기반 Avro Bytes Decoder

이 Avro bytes decoder는 먼저 입력 메시지 바이트에서 subject와 id를 추출한 다음, 그것들을 사용해서 바이트에서 Avro 레코드를 디코딩하는 데 사용할 Avro 스키마를 찾아요. 자세한 내용은 schema repo를 참고하세요. Avro 스키마를 보관하려면 schema repo 같은 HTTP 서비스가 필요합니다. 메시지 생산자 측에서 스키마를 등록하는 방법은 org.apache.druid.data.input.avro.AvroStreamInputFormatTest#testParse()를 참고하세요.

| Field | Type | Description | Required | | type | String | 값을 schema_repo로 설정 | no | | subjectAndIdConverter | JSON Object | 메시지 바이트에서 subject와 id를 추출하는 방법 지정 | yes | | schemaRepository | JSON Object | subject와 id에서 Avro 스키마를 찾는 방법 지정 | yes |

Avro-1124 Subject And Id Converter

이 섹션은 schema_repo Avro bytes decoder의 subjectAndIdConverter 객체 형식을 설명합니다.

| Field | Type | Description | Required | | type | String | 값을 avro_1124로 설정 | no | | topic | String | Kafka 스트림의 토픽 지정 | yes |

Avro-1124 Schema Repository

이 섹션은 schema_repo Avro bytes decoder의 schemaRepository 객체 형식을 설명합니다.

| Field | Type | Description | Required | | type | String | 값을 avro_1124_rest_client로 설정 | no | | url | String | Avro-1124 스키마 저장소의 엔드포인트 URL 지정 | yes |

Confluent Schema Registry 기반 Avro Bytes Decoder

이 Avro bytes decoder는 먼저 입력 메시지 바이트에서 고유 id를 추출한 다음, 그것을 사용해서 바이트에서 Avro 레코드를 디코딩할 때 사용하는 Schema Registry의 스키마를 찾아요. 자세한 내용은 Schema Registry 문서와 저장소를 참고하세요.

| Field | Type | Description | Required | | type | String | 값을 schema_registry로 설정 | no | | url | String | Schema Registry의 URL 엔드포인트 지정 | yes | | capacity | Integer | 캐시의 최대 크기 지정(기본 = Integer.MAX_VALUE) | no | | urls | ARRAY | 여러 Schema Registry 인스턴스의 URL 엔드포인트 지정 | yes (url이 제공되지 않으면) | | config | Json | Schema Registry용으로 구성된 추가 설정을 보내기 위함. DynamicConfigProvider로 제공 가능 | no | | headers | Json | Schema Registry에 헤더를 보내기 위함. DynamicConfigProvider로 제공 가능 | no |

단일 스키마 레지스트리 인스턴스에는 url 필드를 사용하고, 여러 인스턴스에는 urls를 사용하세요.

단일 인스턴스:

..."avroBytesDecoder" : {   "type" : "schema_registry",   "url" : <schema-registry-url>}...

여러 인스턴스:

..."avroBytesDecoder" : {   "type" : "schema_registry",   "urls" : [<schema-registry-url-1>, <schema-registry-url-2>, ...],   "config" : {        "basic.auth.credentials.source": "USER_INFO",        "basic.auth.user.info": "fred:letmein",        "schema.registry.ssl.truststore.location": "/some/secrets/kafka.client.truststore.jks",        "schema.registry.ssl.truststore.password": "<password>",        "schema.registry.ssl.keystore.location": "/some/secrets/kafka.client.keystore.jks",        "schema.registry.ssl.keystore.password": "<password>",        "schema.registry.ssl.key.password": "<password>",        "schema.registry.ssl.key.password",       ...    },   "headers": {       "traceID" : "b29c5de2-0db4-490b-b421",       "timeStamp" : "1577191871865",       "druid.dynamic.config.provider":{            "type":"mapString",             "config":{                 "registry.header.prop.1":"value.1",                  "registry.header.prop.2":"value.2"                 }            }       ...    }}...
Parse exceptions

레코드를 읽을 때 다음 오류는 파싱 예외로 간주되며, maxParseExceptions와 maxSavedParseExceptions 같은 수집 태스크 구성으로 제한·로깅할 수 있어요.

  • 잘못된 구성이나 손상된 레코드(잘못된 스키마 ID)로 인해 스키마를 가져오지 못함
  • Avro 메시지 디코딩 실패

Avro OCF

Avro OCF input format을 로드하려면 Druid Avro 확장(druid-avro-extensions)을 로드하세요.

Druid에서 Avro 타입이 처리되는 방법은 Avro Types 섹션을 참고하세요.

Avro OCF 데이터를 로드하도록 Avro OCF inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 avro_ocf로 설정 | yes | | flattenSpec | JSON Object | Avro 레코드에서 중첩 값을 추출할 flattenSpec 정의. 'path' 표현식만 지원됨('jq'와 'tree'는 사용 불가) | no (기본은 'root' 레벨 속성 자동 발견) | | schema | JSON Object | Avro 레코드를 파싱할 때 사용할 reader 스키마 정의. 여러 버전의 Avro OCF 파일 데이터를 파싱할 때 유용 | no (기본은 OCF 파일에 포함된 writer 스키마로 디코딩) | | binaryAsString | Boolean | 논리적으로 문자열 또는 enum 타입으로 표시되지 않은 bytes parquet 컬럼을 UTF-8 인코딩 문자열로 취급할지 지정 | no (기본 = false) |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "avro_ocf",    "flattenSpec": {      "useFieldDiscovery": true,      "fields": [        {          "type": "path",          "name": "someRecord_subInt",          "expr": "$.someRecord.subInt"        }      ]    },    "schema": {      "namespace": "org.apache.druid.data.input",      "name": "SomeDatum",      "type": "record",      "fields" : [        { "name": "timestamp", "type": "long" },        { "name": "eventType", "type": "string" },        { "name": "id", "type": "long" },        { "name": "someRecord", "type": {          "type": "record", "name": "MySubRecord", "fields": [            { "name": "subInt", "type": "int"},            { "name": "subLong", "type": "long"}          ]        }}]    },    "binaryAsString": false  },  ...}

Protobuf

정보

Protobuf input format을 사용하려면 druid-protobuf-extensions를 확장으로 포함해야 합니다.

Protobuf 데이터를 로드하도록 Protobuf inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | | type | String | 값을 protobuf로 설정 | yes | | flattenSpec | JSON Object | Protobuf 레코드에서 중첩 값을 추출할 flattenSpec 정의. 'path' 표현식만 지원됨('jq'와 'tree'는 사용 불가)에 유의 | no (기본은 'root' 레벨 속성 자동 발견) | | protoBytesDecoder | JSON Object | 바이트를 Protobuf 레코드로 디코딩하는 방법 지정 | yes |

예를 들어:

"ioConfig": {  "inputFormat": {    "type": "protobuf",    "protoBytesDecoder": {      "type": "file",      "descriptor": "file:///tmp/metrics.desc",      "protoMessageType": "Metrics"    }    "flattenSpec": {      "useFieldDiscovery": true,      "fields": [        {          "type": "path",          "name": "someRecord_subInt",          "expr": "$.someRecord.subInt"        }      ]    }  },  ...}

Kafka

kafka input format은 Kafka 페이로드 값 내용에 더해 Kafka 메타데이터 필드를 파싱할 수 있게 해 줍니다. Apache Kafka에서 수집할 때만 사용해야 해요.

kafka input format은 페이로드 파싱 input format을 감싸고, 그 출력에 Kafka 이벤트 타임스탬프, 토픽 이름, 이벤트 헤더, 그리고 그 자체도 사용 가능한 어떤 input format으로든 파싱할 수 있는 key 필드를 추가합니다.

페이로드의 컬럼 이름과 메타데이터에서 생성된 컬럼 이름 사이에 충돌이 있으면 페이로드가 우선합니다. 이렇게 함으로써 Kafka 수집을 Kafka input format으로 업그레이드(기존 input format을 가져와 valueFormat으로 설정)해도 페이로드 데이터가 손실되지 않습니다.

Kafka inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | Default | | type | String | 값을 kafka로 설정 | yes | | | valueFormat | InputFormat | Kafka 값 페이로드를 파싱할 input format | yes | | | timestampColumnName | String | Kafka 타임스탬프용 컬럼 이름 | no | kafka.timestamp | | topicColumnName | String | Kafka 토픽용 컬럼 이름. 여러 토픽에서 같은 datasource로 데이터를 수집할 때 유용 | no | kafka.topic | | headerColumnPrefix | String | 모든 헤더 컬럼의 커스텀 접두사 | no | kafka.header | | headerFormat | Object | Kafka 헤더를 파싱하는 방법 지정. String 타입 지원. Kafka 헤더 값은 바이트이므로 파서가 그것들을 UTF-8 인코딩 문자열로 디코딩. 이 동작을 바꾸려면 인코딩 스타일을 기반으로 나만의 파서를 구현. KafkaStringHeaderFormat의 encoding 타입을 커스텀 구현에 맞게 변경. 지원 인코딩 형식은 Header format 참고 | no | | | keyFormat | InputFormat | Kafka key를 파싱할 input format. inputFormat 필드의 첫 번째 항목만 처리. key 값이 단순 문자열이면 tsv format으로 파싱할 수 있음. tsv, csv, regex 형식에서는 유효한 input format을 만들기 위해 columns 배열을 제공해야 함에 유의. 첫 번째 것만 사용되고, 그 이름은 keyColumnName을 대신 사용하므로 무시됨 | no | | | keyColumnName | String | Kafka key용 컬럼 이름 | no | kafka.key |

Header format

headerFormat은 다음 인코딩 형식을 지원합니다.

  • ISO-8859-1: ISO Latin Alphabet No. 1, 즉 ISO-LATIN-1.
  • US-ASCII: 일곱 비트 ASCII. ISO646-US라고도 함. Unicode 문자 집합의 Basic Latin 블록.
  • UTF-8: 여덟 비트 UCS Transformation Format.
  • UTF-16: 십육 비트 UCS Transformation Format, 선택적 byte-order mark로 바이트 순서 식별.
  • UTF-16BE: 십육 비트 UCS Transformation Format, big-endian 바이트 순서.
  • UTF-16LE: 십육 비트 UCS Transformation Format, little-endian 바이트 순서.
  • headerColumnPrefix: 페이로드 컬럼과의 충돌을 피하기 위해 Kafka 헤더에 접두사 제공. 기본값은 kafka.header..

예시

input format으로 { "type": "json" }을 사용하면 페이로드 값만 파싱돼요. 페이로드에 더해 Kafka 메타데이터까지 파싱하려면 kafka input format을 사용하세요.

예를 들어 개발 환경에서 편집을 나타내는 Kafka 메시지의 다음 구조를 고려해 보세요.

  • Kafka 타임스탬프: 1680795276351
  • Kafka 토픽: wiki-edits
  • Kafka 헤더: - env=development - zone=z1
  • Kafka key: wiki-edit
  • Kafka 페이로드 값: {"channel":"#sv.wikipedia","timestamp":"2016-06-27T00:00:11.080Z","page":"Salo Toraut","delta":31,"namespace":"Main"}

다음과 같이 구성하면 됩니다.

"ioConfig": {  "inputFormat": {    "type": "kafka",    "valueFormat": {      "type": "json"    },    "timestampColumnName": "kafka.timestamp",    "topicColumnName": "kafka.topic",    "headerFormat": {      "type": "string",      "encoding": "UTF-8"    },    "headerColumnPrefix": "kafka.header.",    "keyFormat": {      "type": "tsv",      "findColumnsFromHeader": false,      "columns": ["x"]    },    "keyColumnName": "kafka.key",  }}

예시 메시지를 다음과 같이 파싱합니다.

{  "channel": "#sv.wikipedia",  "timestamp": "2016-06-27T00:00:11.080Z",  "page": "Salo Toraut",  "delta": 31,  "namespace": "Main",  "kafka.timestamp": 1680795276351,  "kafka.topic": "wiki-edits",  "kafka.header.env": "development",  "kafka.header.zone": "z1",  "kafka.key": "wiki-edit"}

kafka.timestamp를 Druid의 기본 타임스탬프(__time)로 사용하고 싶다면, timestampSpec의 column 값으로 지정하세요.

"timestampSpec": {  "column": "kafka.timestamp",  "format": "millis"}

마찬가지로 Kafka 헤더에서 추출한 타임스탬프를 사용하고 싶다면:

"timestampSpec": {  "column": "kafka.header.myTimestampHeader",  "format": "millis"}

마지막으로 이 Kafka 메타데이터 컬럼들을 dimensionsSpec에 추가하거나, dimensionsSpec이 컬럼을 자동 감지하도록 설정하세요.

다음 supervisor spec은 Kafka 헤더·key·타임스탬프·토픽을 Druid dimension으로 수집하는 방법을 보여줘요.

예시 보기

{  "type": "kafka",  "spec": {    "ioConfig": {      "type": "kafka",      "consumerProperties": {        "bootstrap.servers": "localhost:9092"      },      "topic": "wiki-edits",      "inputFormat": {        "type": "kafka",        "valueFormat": {          "type": "json"        },        "headerFormat": {          "type": "string"        },        "keyFormat": {          "type": "tsv",          "findColumnsFromHeader": false,          "columns": ["x"]        }      },      "useEarliestOffset": true    },    "dataSchema": {      "dataSource": "wikiticker",      "timestampSpec": {        "column": "timestamp",        "format": "posix"      },      "dimensionsSpec":  "dimensionsSpec": {        "useSchemaDiscovery": true,        "includeAllDimensions": true      },      "granularitySpec": {        "queryGranularity": "none",        "rollup": false,        "segmentGranularity": "day"      }    },    "tuningConfig": {      "type": "kafka"    }  }}

Druid가 데이터를 수집한 뒤에는 다음과 같이 Kafka 메타데이터 컬럼을 쿼리할 수 있어요.

SELECT  "kafka.header.env",  "kafka.key",  "kafka.timestamp",  "kafka.topic"FROM "wikiticker"

이 쿼리는 다음을 반환합니다.

| kafka.header.env | kafka.key | kafka.timestamp | kafka.topic | | development | wiki-edit | 1680795276351 | wiki-edits |

Kinesis

kinesis input format은 Kinesis 페이로드 값 내용에 더해 Kinesis 메타데이터 필드를 파싱할 수 있게 해 줍니다. Kinesis에서 수집할 때만 사용해야 해요.

kinesis input format은 페이로드 파싱 input format을 감싸고, 그 출력에 Kinesis 이벤트 타임스탬프와 파티션 키, 즉 Kinesis 레코드의 ApproximateArrivalTimestamp 와 PartitionKey 필드를 추가합니다.

페이로드의 컬럼 이름과 메타데이터에서 생성된 컬럼 이름 사이에 충돌이 있으면 페이로드가 우선합니다. 이렇게 함으로써 Kinesis 수집을 Kinesis input format으로 업그레이드(기존 input format을 가져와 valueFormat으로 설정)해도 페이로드 데이터가 손실되지 않습니다.

Kinesis inputFormat을 다음과 같이 구성하세요.

| Field | Type | Description | Required | Default | | type | String | 값을 kinesis로 설정 | yes | | | valueFormat | InputFormat | Kinesis 값 페이로드를 파싱할 input format | yes | | | partitionKeyColumnName | String | Kinesis 파티션 키용 컬럼 이름. 여러 파티션에서 같은 datasource로 데이터를 수집할 때 유용 | no | kinesis.partitionKey | | timestampColumnName | String | Kinesis 타임스탬프용 컬럼 이름 | no | kinesis.timestamp |

예시

input format으로 { "type": "json" }을 사용하면 페이로드 값만 파싱돼요. 페이로드에 더해 Kinesis 메타데이터까지 파싱하려면 kinesis input format을 사용하세요.

예를 들어 개발 환경에서 편집을 나타내는 Kinesis 레코드의 다음 구조를 고려해 보세요.

  • Kinesis 타임스탬프: 1680795276351
  • Kinesis 파티션 키: partition-1
  • Kinesis 페이로드 값: {"channel":"#sv.wikipedia","timestamp":"2016-06-27T00:00:11.080Z","page":"Salo Toraut","delta":31,"namespace":"Main"}

다음과 같이 구성하면 됩니다.

{  "ioConfig": {    "inputFormat": {      "type": "kinesis",      "valueFormat": {        "type": "json"      },      "timestampColumnName": "kinesis.timestamp",      "partitionKeyColumnName": "kinesis.partitionKey"    }  }}

예시 레코드를 다음과 같이 파싱합니다.

{  "channel": "#sv.wikipedia",  "timestamp": "2016-06-27T00:00:11.080Z",  "page": "Salo Toraut",  "delta": 31,  "namespace": "Main",  "kinesis.timestamp": 1680795276351,  "kinesis.partitionKey": "partition-1"}

kinesis.timestamp를 Druid의 기본 타임스탬프(__time)로 사용하고 싶다면, timestampSpec의 column 값으로 지정하세요.

"timestampSpec": {  "column": "kinesis.timestamp",  "format": "millis"}

마지막으로 이 Kinesis 메타데이터 컬럼들을 dimensionsSpec에 추가하거나, dimensionsSpec이 컬럼을 자동 감지하도록 설정하세요.

다음 supervisor spec은 Kinesis 타임스탬프와 파티션 키를 Druid dimension으로 수집하는 방법을 보여줘요.

예시 보기

{  "type": "kinesis",  "spec": {    "ioConfig": {      "type": "kinesis",      "consumerProperties": {        "bootstrap.servers": "localhost:9092"      },      "topic": "wiki-edits",      "inputFormat": {        "type": "kinesis",        "valueFormat": {          "type": "json"        },        "headerFormat": {          "type": "string"        },        "keyFormat": {          "type": "tsv",          "findColumnsFromHeader": false,          "columns": ["x"]        }      },      "useEarliestOffset": true    },    "dataSchema": {      "dataSource": "wikiticker",      "timestampSpec": {        "column": "timestamp",        "format": "posix"      },      "dimensionsSpec": {        "useSchemaDiscovery": true,        "includeAllDimensions": true      },      "granularitySpec": {        "queryGranularity": "none",        "rollup": false,        "segmentGranularity": "day"      }    },    "tuningConfig": {      "type": "kinesis"    }  }}

Druid가 데이터를 수집한 뒤에는 다음과 같이 Kinesis 메타데이터 컬럼을 쿼리할 수 있어요.

SELECT  "kinesis.timestamp",  "kinesis.partitionKey"FROM "wikiticker"

이 쿼리는 다음을 반환합니다.

| kinesis.timestamp | kinesis.topic | | 1680795276351 | partition-1 |

FlattenSpec

중첩 데이터를 평탄화하는 데 Druid 중첩 컬럼 기능의 대안으로, 그리고 그 기능이 지원하지 않는 중첩 input format에서 flattenSpec 객체를 사용할 수 있어요. inputFormat 객체 안에 있는 객체입니다.

중첩 데이터를 COMPLEX<json> 데이터 타입의 Apache Druid 컬럼으로 수집·저장하는 방법은 Nested columns를 참고하세요.

flattenSpec을 다음과 같이 구성하세요.

| Field | Description | Default | | useFieldDiscovery | true면 모든 root-level 필드를 timestampSpec, transformSpec, dimensionsSpec, metricsSpec이 사용할 수 있는 사용 가능한 필드로 해석. false면 명시적으로 지정된 필드(fields 참고)만 사용 가능 | true | | fields | 관심 있는 필드와 그것들에 접근하는 방법 지정. 자세한 내용은 Field flattening specifications 참고 | [] |

예를 들어:

"flattenSpec": {  "useFieldDiscovery": true,  "fields": [    { "name": "baz", "type": "root" },    { "name": "foo_bar", "type": "path", "expr": "$.foo.bar" },    { "name": "foo_other_bar", "type": "tree", "nodes": ["foo", "other", "bar"] },    { "name": "first_food", "type": "jq", "expr": ".thing.food[1]" }  ]}

Druid가 입력 데이터 레코드를 읽은 뒤, 다른 스펙(timestampSpec, transformSpec, dimensionsSpec, metricsSpec)을 적용하기 전에 flattenSpec을 적용해요. 이렇게 하면 예를 들어 평탄화된 데이터에서 타임스탬프를 추출하고, 변환·dimension 목록·메트릭 생성에서 평탄화된 데이터를 참조할 수 있습니다.

평탄화는 avro, json, orc, parquet를 포함해 중첩을 지원하는 데이터 형식에서만 지원됩니다.

Field flattening specifications

fields 목록의 각 항목은 다음 컴포넌트를 가질 수 있어요.

| Field | Description | Default | | type | 옵션은 다음과 같음: - root: 레코드의 root 레벨에 있는 필드 참조. useFieldDiscovery가 false일 때만 실제로 유용. - path: JsonPath 표기법으로 필드 참조. avro, json, orc, parquet를 포함해 중첩을 제공하는 대부분의 데이터 형식에서 지원. - jq: jackson-jq 표기법으로 필드 참조. json 형식에서만 지원. - tree: 레코드의 root 레벨에서 중첩 필드 참조. 단순한 계층적 가져오기가 필요하다면 path나 jq보다 유용하고 효율적. json 형식에서만 지원 | none (필수) | | name | 평탄화 후 필드 이름. 이 이름은 timestampSpec, transformSpec, dimensionsSpec, metricsSpec이 참조할 수 있음 | none (필수) | | expr | 평탄화 중 필드에 접근하기 위한 표현식. path 타입이면 JsonPath여야 하고, jq 타입이면 jackson-jq 표기법이어야 함. 다른 타입에서는 이 파라미터가 무시됨 | none (path와 jq 타입에서 필수) | | nodes | tree 전용. 평탄화 중 필드에 접근하기 위한 다중 표현식 필드로, 읽을 필드 이름의 계층을 나타냄. 다른 타입에서는 이 파라미터를 제공하면 안 됨 | none (tree 타입에서 필수) |

평탄화에 관한 참고 사항

  • 편의를 위해 root-level 필드를 정의할 때 JSON 객체 대신 문자열로 필드 이름만 정의할 수 있어요. 예를 들어 {"name": "baz", "type": "root"}는 "baz"와 동등합니다.
  • useFieldDiscovery를 활성화하면 Druid가 지원하는 데이터 타입에 해당하는 root 레벨의 "단순한" 필드만 자동 감지해요. 여기에는 문자열, 숫자, 문자열·숫자 목록이 포함됩니다. 다른 타입은 자동 감지되지 않으며 fields 목록에 명시적으로 지정해야 해요.
  • 중복 필드 name은 허용되지 않습니다. 예외가 던져질 거예요.
  • useFieldDiscovery가 활성화되면, 이미 fields 목록에 정의된 것과 같은 이름의 발견된 필드는 두 번 추가되는 대신 건너뜁니다.
  • JSONPath 평가기는 path-타입 표현식을 테스트하는 데 유용해요.
  • jackson-jq는 전체 jq 문법의 부분 집합을 지원합니다. 자세한 내용은 jackson-jq 문서를 참고하세요.
  • JsonPath는 많은 함수를 지원하지만, 그 모든 함수를 Druid가 현재 지원하는 것은 아니에요. 다음 매트릭스는 현재 지원되는 JsonPath 함수와 해당 데이터 형식을 보여줍니다. 이 함수들의 출력 데이터 타입에도 유의하세요.

| Function | Description | Output type | json | orc | avro | parquet | | min() | 숫자 배열의 최소값 제공 | Double | ✓ | ✓ | ✓ | ✓ | | max() | 숫자 배열의 최대값 제공 | Double | ✓ | ✓ | ✓ | ✓ | | avg() | 숫자 배열의 평균값 제공 | Double | ✓ | ✓ | ✓ | ✓ | | stddev() | 숫자 배열의 표준 편차 값 제공 | Double | ✓ | ✓ | ✓ | ✓ | | length() | 배열의 길이 제공 | Integer | ✓ | ✓ | ✓ | ✓ | | sum() | 숫자 배열의 합계 값 제공 | Double | ✓ | ✓ | ✓ | ✓ | | concat(X) | 새 항목과 함께 path 출력의 연결된 버전 제공 | 입력과 동일 | ✓ | ✗ | ✗ | ✗ | | append(X) | json path 출력 배열에 항목 추가 | 입력과 동일 | ✓ | ✗ | ✗ | ✗ | | keys() | 속성 키 제공 (종단 물결표 ~의 대안) | Set | ✗ | ✗ | ✗ | ✗ |

더 알아보기 (Learn more)