Kafka Connect 사용자 가이드

Kafka Connect 사용자 가이드 (User Guide)

이 페이지는 Kafka Connect를 실제로 구성하고 실행하고 관리하는 방법을 자세히 다뤄요. standalone(단독) 모드와 distributed(분산) 모드의 차이, 커넥터 구성, 트랜스포메이션(변환), REST API, 장애 처리, 정확히 한 번(exactly-once) 지원 같은 운영에 꼭 필요한 내용이 촘촘하게 담겨 있어요.

출처: 문서

본문

quickstart는 standalone 버전의 Kafka Connect를 실행하는 간단한 예시를 제공합니다. 이 섹션에서는 Kafka Connect를 구성·실행·관리하는 방법을 더 자세히 설명합니다.

Kafka Connect 실행 (Running Kafka Connect)

Kafka Connect는 현재 standalone(단일 프로세스)과 distributed(분산) 두 가지 실행 모드를 지원합니다.

standalone 모드에서는 모든 작업이 단일 프로세스에서 수행됩니다. 이 구성은 설정하고 시작하기가 더 간단하며, 워커(worker)가 하나만 의미 있는 상황(예: 로그 파일 수집)에서 유용할 수 있지만, 장애 허용(fault tolerance) 같은 Kafka Connect의 일부 기능은 누릴 수 없습니다. 다음 명령으로 standalone 프로세스를 시작할 수 있습니다.

$ bin/connect-standalone.sh config/connect-standalone.properties [connector1.properties connector2.json …]

첫 번째 파라미터는 워커의 구성입니다. 여기에는 Kafka 연결 파라미터, 직렬화 형식, 오프셋 커밋 빈도 같은 설정이 포함됩니다. 제공된 예시는 config/server.properties가 제공하는 기본 구성으로 실행 중인 로컬 클러스터에서 잘 작동합니다. 다른 구성이나 프로덕션 배포에 사용하려면 조정이 필요합니다. 모든 워커(standalone과 distributed 모두)는 몇 가지 구성이 필요합니다.

  • bootstrap.servers - Kafka에 대한 연결을 부트스트랩하는 데 사용되는 Kafka 서버 목록.
  • key.converter - Kafka Connect 형식과 Kafka에 기록되는 직렬화된 형식 사이를 변환하는 데 사용되는 컨버터 클래스. Kafka에서 쓰거나 읽는 메시지의 키 형식을 제어하며, 커넥터와 독립적이므로 어떤 커넥터든 어떤 직렬화 형식과도 함께 작동하게 해줍니다. 일반적인 형식의 예로 JSON과 Avro가 있습니다.
  • value.converter - Kafka Connect 형식과 Kafka에 기록되는 직렬화된 형식 사이를 변환하는 데 사용되는 컨버터 클래스. Kafka에서 쓰거나 읽는 메시지의 값 형식을 제어하며, 커넥터와 독립적이므로 어떤 커넥터든 어떤 직렬화 형식과도 함께 작동하게 해줍니다. 일반적인 형식의 예로 JSON과 Avro가 있습니다.
  • plugin.path(기본값 null) - Connect 플러그인(커넥터, 컨버터, 트랜스포메이션)이 포함된 경로 목록. quickstart를 실행하기 전에 사용자는 connect-file-4.3.1.jar에 패키징된 예시 FileStreamSourceConnectorFileStreamSinkConnector가 포함된 절대 경로를 추가해야 합니다. 이 커넥터들은 기본적으로 Connect 워커의 CLASSPATH나 plugin.path에 포함되지 않기 때문입니다.

standalone 모드에 특화된 중요한 구성 옵션은 다음과 같습니다.

  • offset.storage.file.filename - 소스 커넥터 오프셋을 저장하는 파일.

여기서 구성된 파라미터는 Kafka Connect가 구성·오프셋·상태 토픽에 접근하는 데 사용하는 프로듀서와 컨슈머를 위한 것입니다. Kafka 소스 태스크가 사용하는 프로듀서와 Kafka 싱크 태스크가 사용하는 컨슈머의 구성에는 동일한 파라미터를 사용할 수 있지만, 각각 producer.consumer. 접두사를 붙여야 합니다. 워커 구성에서 접두사 없이 상속되는 유일한 Kafka 클라이언트 파라미터는 bootstrap.servers이며, 대부분의 경우 같은 클러스터가 모든 용도에 사용되므로 이는 대개 충분합니다. 주목할 예외는 보안 클러스터로, 연결을 허용하려면 추가 파라미터가 필요합니다. 이 파라미터는 관리 접근, Kafka 소스, Kafka 싱크 등 워커 구성에서 최대 세 번 설정해야 합니다.

클라이언트 구성 오버라이드는 Kafka 소스나 Kafka 싱크에 대해 각각 producer.override.consumer.override. 접두사를 사용해 커넥터별로 개별 구성할 수 있습니다. 이러한 오버라이드는 커넥터의 나머지 구성 속성과 함께 포함됩니다.

나머지 파라미터는 커넥터 구성 파일입니다. 각 파일은 Java Properties 파일이거나, POST /connectors 엔드포인트나 PUT /connectors/{name}/config 엔드포인트의 요청 본문과 동일한 구조의 객체를 담은 JSON 파일일 수 있습니다. 원하는 만큼 많이 포함할 수 있지만, 모두 동일한 프로세스 내에서(다른 스레드에서) 실행됩니다. 커맨드 라인에서 커넥터 구성 파일을 지정하지 않고, standalone 워커가 시작된 후 런타임에 REST API를 사용해 커넥터를 생성할 수도 있습니다.

distributed 모드는 작업의 자동 균형 조정을 처리하고, 동적으로 확장(또는 축소)할 수 있으며, 활성 태스크와 구성·오프셋 커밋 데이터 모두에 장애 허용을 제공합니다. 실행은 standalone 모드와 매우 유사합니다.

$ bin/connect-distributed.sh config/connect-distributed.properties

차이는 시작되는 클래스와, Kafka Connect 프로세스가 구성 저장 위치, 작업 할당 방식, 오프셋·태스크 상태 저장 위치를 결정하는 방식을 바꾸는 구성 파라미터에 있습니다. distributed 모드에서 Kafka Connect는 오프셋, 구성, 태스크 상태를 Kafka 토픽에 저장합니다. 원하는 파티션 수와 복제 팩터를 얻으려면 오프셋, 구성, 상태용 토픽을 수동으로 생성하는 것이 권장됩니다. Kafka Connect 시작 시 이 토픽들이 아직 생성되지 않았다면 기본 파티션 수와 복제 팩터로 자동 생성되는데, 이는 그 용도에 가장 적합하지 않을 수 있습니다.

특히 다음 구성 파라미터는 위에서 언급한 공통 설정에 추가로, 클러스터를 시작하기 전에 설정하는 것이 중요합니다.

  • group.id - 클러스터의 고유 이름. Connect 클러스터 그룹을 형성하는 데 사용됩니다. 컨슈머 그룹 ID와 충돌하지 않아야 합니다.
  • config.storage.topic - 커넥터와 태스크 구성을 저장하는 데 사용할 토픽 이름. 이 토픽은 단일 파티션이어야 하고, 복제되고, 컴팩션용으로 구성되어야 합니다.
  • offset.storage.topic - 오프셋을 저장하는 데 사용할 토픽 이름. 이 토픽은 파티션이 많아야 하고, 복제되고, 컴팩션용으로 구성되어야 합니다.
  • status.storage.topic - 상태를 저장하는 데 사용할 토픽 이름. 이 토픽은 여러 파티션을 가질 수 있고, 복제되고, 컴팩션용으로 구성되어야 합니다.

distributed 모드에서는 커넥터 구성이 커맨드 라인에서 전달되지 않습니다. 대신 아래에 설명된 REST API를 사용해 커넥터를 생성·수정·삭제합니다.

커넥터 구성 (Configuring Connectors)

커넥터 구성은 단순한 키-값 매핑입니다. standalone과 distributed 모드 모두에서 커넥터를 생성(또는 수정)하는 REST 요청의 JSON 페이로드에 포함됩니다. standalone 모드에서는 properties 파일로 정의해 커맨드 라인에서 Connect 프로세스에 전달할 수도 있습니다.

대부분의 구성은 커넥터마다 다르므로 여기서 모두 설명할 수 없습니다. 하지만 몇 가지 공통 옵션이 있습니다.

  • name - 커넥터의 고유 이름. 같은 이름으로 다시 등록하려 하면 실패합니다.
  • connector.class - 커넥터의 Java 클래스.
  • tasks.max - 이 커넥터에 대해 생성해야 하는 최대 태스크 수. 커넥터가 이 수준의 병렬성을 달성하지 못하면 더 적은 태스크를 생성할 수 있습니다.
  • key.converter(선택) - 워커가 설정한 기본 키 컨버터를 오버라이드.
  • value.converter(선택) - 워커가 설정한 기본 값 컨버터를 오버라이드.

connector.class 구성은 여러 형식을 지원합니다. 이 커넥터 클래스의 전체 이름 또는 별칭(alias)입니다. 커넥터가 org.apache.kafka.connect.file.FileStreamSinkConnector라면 전체 이름을 지정하거나, 구성을 조금 더 짧게 하기 위해 FileStreamSink 또는 FileStreamSinkConnector를 사용할 수 있습니다.

싱크 커넥터는 입력을 제어하는 몇 가지 추가 옵션도 있습니다. 각 싱크 커넥터는 다음 중 하나를 설정해야 합니다.

  • topics - 이 커넥터의 입력으로 사용할 토픽의 쉼표로 구분된 목록.
  • topics.regex - 이 커넥터의 입력으로 사용할 토픽의 Java 정규식.

다른 옵션은 커넥터의 문서를 참고해야 합니다.

트랜스포메이션 (Transformations)

커넥터는 메시지를 한 번에 하나씩 가볍게 수정하는 트랜스포메이션으로 구성할 수 있습니다. 데이터 다듬기(data massaging)와 이벤트 라우팅에 편리할 수 있습니다.

커넥터 구성에서 트랜스포메이션 체인(chain)을 지정할 수 있습니다.

  • transforms - 트랜스포메이션의 별칭 목록. 트랜스포메이션이 적용될 순서를 지정합니다.
  • transforms.$alias.type - 트랜스포메이션의 정규화된 클래스 이름.
  • transforms.$alias.$transformationSpecificConfig - 트랜스포메이션의 구성 속성.

예를 들어 내장 파일 소스 커넥터를 사용해 정적 필드를 추가하는 트랜스포메이션을 적용해 봅시다.

예제 전체에서 스키마 없는(schemaless) JSON 데이터 형식을 사용하겠습니다. 스키마 없는 형식을 사용하기 위해 connect-standalone.properties의 다음 두 줄을 true에서 false로 바꿨습니다.

key.converter.schemas.enable
value.converter.schemas.enable

파일 소스 커넥터는 각 줄을 String으로 읽습니다. 각 줄을 Map으로 감싼 다음, 이벤트 출처를 식별하는 두 번째 필드를 추가하겠습니다. 이를 위해 두 개의 트랜스포메이션을 사용합니다.

  • HoistField - 입력 줄을 Map 안에 배치.
  • InsertField - 정적 필드 추가. 이 예제에서는 레코드가 파일 커넥터에서 왔다는 것을 표시.

트랜스포메이션을 추가한 후 connect-file-source.properties 파일은 다음과 같습니다.

name=local-file-source
connector.class=FileStreamSource
tasks.max=1
file=test.txt
topic=connect-test
transforms=MakeMap, InsertSource
transforms.MakeMap.type=org.apache.kafka.connect.transforms.HoistField$Value
transforms.MakeMap.field=line
transforms.InsertSource.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.InsertSource.static.field=data_source
transforms.InsertSource.static.value=test-file-source

transforms로 시작하는 모든 줄은 트랜스포메이션을 위해 추가된 것입니다. 우리가 만든 두 트랜스포메이션 "InsertSource"와 "MakeMap"은 우리가 트랜스포메이션에 붙인 별칭입니다. 트랜스포메이션 타입은 아래에서 볼 수 있는 내장 트랜스포메이션 목록을 기반으로 합니다. 각 트랜스포메이션 타입에는 추가 구성이 있습니다. HoistField는 파일의 원래 String을 담을 Map의 필드 이름인 "field"라는 구성이 필요합니다. InsertField 트랜스포메이션은 추가하는 필드 이름과 값을 지정하게 해줍니다.

샘플 파일에서 트랜스포메이션 없이 파일 소스 커넥터를 실행하고 kafka-console-consumer.sh로 읽었을 때 결과는 다음과 같았습니다.

"foo"
"bar"
"hello world"

그다음 새 파일 커넥터를 만들되, 이번에는 구성 파일에 트랜스포메이션을 추가했습니다. 이번 결과는 다음과 같습니다.

{"line":"foo","data_source":"test-file-source"}
{"line":"bar","data_source":"test-file-source"}
{"line":"hello world","data_source":"test-file-source"}

읽은 줄이 이제 JSON map의 일부가 되었고, 우리가 지정한 정적 값이 담긴 추가 필드가 생긴 것을 볼 수 있습니다. 이것은 트랜스포메이션으로 할 수 있는 일의 한 예일 뿐입니다.

포함된 트랜스포메이션 (Included transformations)

Kafka Connect에는 널리 적용 가능한 여러 데이터·라우팅 트랜스포메이션이 포함되어 있습니다.

  • Cast - 필드 또는 전체 키·값을 특정 타입으로 캐스팅.
  • DropHeaders - 이름으로 헤더 제거.
  • ExtractField - Struct와 Map에서 특정 필드 추출해 결과에 그 필드만 포함.
  • Filter - 이후 모든 처리에서 메시지 제거. 프레디케이트(predicate)와 함께 사용해 특정 메시지를 선택적으로 걸러냄.
  • Flatten - 중첩 데이터 구조 펼치기.
  • HeaderFrom - 키나 값의 필드를 레코드 헤더로 복사하거나 이동.
  • HoistField - 전체 이벤트를 Struct나 Map 안의 단일 필드로 감싸기.
  • InsertField - 정적 데이터 또는 레코드 메타데이터를 사용해 필드 추가.
  • InsertHeader - 정적 데이터를 사용해 헤더 추가.
  • MaskField - 필드를 타입에 맞는 유효한 null 값(0, 빈 문자열 등) 또는 커스텀 치환 값으로 대체.
  • RegexRouter - 원래 토픽, 치환 문자열, 정규식을 기반으로 레코드 토픽 수정.
  • ReplaceField - 필드 필터링 또는 이름 변경.
  • SetSchemaMetadata - 스키마 이름이나 버전 수정.
  • TimestampConverter - 타임스탬프를 다른 형식 간에 변환.
  • TimestampRouter - 원래 토픽과 타임스탬프를 기반으로 레코드 토픽 수정. 타임스탬프 기반으로 다른 테이블이나 인덱스에 써야 하는 싱크를 사용할 때 유용.
  • ValueToKey - 레코드 값의 필드 하위 집합으로 형성된 새 키로 레코드 키 교체.

각 트랜스포메이션의 구성 방법 세부 사항은 아래에 나열되어 있습니다.

org.apache.kafka.connect.transforms.Cast

필드 또는 전체 키·값을 특정 타입으로 캐스팅합니다. 예: 정수 필드를 더 작은 폭으로 강제. 정수·부동소수·불리언·문자열에서 다른 타입으로 캐스팅하고, 바이너리를 문자열(base64 인코딩)로 캐스팅합니다.

키(org.apache.kafka.connect.transforms.Cast$Key) 또는 값(org.apache.kafka.connect.transforms.Cast$Value)용으로 설계된 구체적인 트랜스포메이션 타입을 사용하세요.

  • spec - Map 또는 Struct의 필드를 캐스팅하기 위한 field1:type,field2:type 형태의 필드·타입 목록. 전체 값을 캐스팅할 단일 타입. 유효한 타입은 int8, int16, int32, int64, float32, float64, boolean, string. 바이너리 필드는 문자열로만 캐스팅할 수 있습니다. Type: list, Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. true이면 기본값이 사용되고, 아니면 null이 사용됩니다. Type: boolean, Default: true, Importance: medium.
org.apache.kafka.connect.transforms.DropHeaders

각 레코드에서 하나 이상의 헤더를 제거합니다.

  • headers - 제거할 헤더의 이름. Type: list, Importance: high.
org.apache.kafka.connect.transforms.ExtractField

스키마가 있으면 Struct, 스키마 없는 데이터면 Map에서 지정된 필드를 추출합니다. null 값은 수정 없이 그대로 통과됩니다.

키(ExtractField$Key) 또는 값(ExtractField$Value)용 구체적 타입을 사용하세요.

  • field - 추출할 필드 이름. Type: string, Default: (없음), Importance: medium.
  • field.syntax.version - 필드 접근 문법 버전. V1이면 필드 경로가 struct/map의 루트 레벨 요소만 접근할 수 있습니다. V2면 중첩 요소 접근을 지원하며 점 표기법을 사용합니다. 필드 이름에 점이 이미 포함되어 있다면 점이 포함된 필드 이름을 감싸기 위해 백틱 쌍을 사용할 수 있습니다. 예: "foo.bar"라는 필드에서 하위 필드 baz에 접근하려면 "foo.bar.baz" 형식을 사용. Type: string, Default: V1, Valid Values: [V1, V2](대소문자 무시), Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
org.apache.kafka.connect.transforms.Filter

모든 레코드를 버려 체인에서 이후 트랜스포메이션에서 걸러냅니다. 특정 Predicate와 일치(또는 비일치)하는 레코드를 걸러내기 위해 조건부로 사용하기 위한 것입니다.

org.apache.kafka.connect.transforms.Flatten

중첩 데이터 구조를 펼쳐 각 레벨의 필드 이름을 구성 가능한 구분 문자로 연결해 각 필드의 이름을 생성합니다. 스키마 있으면 Struct, 스키마 없는 데이터면 Map에 적용됩니다. 배열 필드와 그 내용은 수정되지 않습니다. 기본 구분 문자는 .입니다.

키(Flatten$Key) 또는 값(Flatten$Value)용 구체적 타입을 사용하세요.

  • delimiter - 출력 레코드의 필드 이름을 생성할 때 입력 레코드의 필드 이름 사이에 삽입할 구분 문자. Type: string, Default: ., Importance: medium.
org.apache.kafka.connect.transforms.HeaderFrom

레코드의 키/값의 필드를 그 레코드의 헤더로 이동하거나 복사합니다. fieldsheaders의 대응 요소가 함께 필드와 그것이 이동·복사될 헤더를 식별합니다. 키(HeaderFrom$Key) 또는 값(HeaderFrom$Value)용 구체적 타입을 사용하세요.

  • fields - 값을 헤더로 복사·이동할 레코드의 필드 이름. Type: list, Valid Values: non-empty list, Importance: high.
  • headers - fields 구성 속성에 나열된 필드 이름과 같은 순서의 헤더 이름. Type: list, Valid Values: non-empty list, Importance: high.
  • operation - 필드를 헤더로 이동(move, 키/값에서 제거)할지 복사(copy, 키/값에 유지)할지. Type: string, Valid Values: [move, copy], Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
org.apache.kafka.connect.transforms.HoistField

스키마 있으면 Struct, 스키마 없는 데이터면 Map에서 지정된 필드 이름을 사용해 데이터를 감쌉니다.

키(HoistField$Key) 또는 값(HoistField$Value)용 구체적 타입을 사용하세요.

  • field - 결과 Struct 또는 Map에서 생성될 단일 필드의 필드 이름. Type: string, Default: (없음), Importance: medium.
org.apache.kafka.connect.transforms.InsertField

레코드 메타데이터의 속성 또는 구성된 정적 값을 사용해 필드를 삽입합니다.

키(InsertField$Key) 또는 값(InsertField$Value)용 구체적 타입을 사용하세요.

  • offset.field - Kafka 오프셋용 필드 이름. 싱크 커넥터에만 적용. !를 붙여 필수 필드로, ?를 붙여 선택(기본)으로. Type: string, Default: null, Importance: medium.
  • partition.field - Kafka 파티션용 필드 이름. ! 필수 / ? 선택(기본). Type: string, Default: null, Importance: medium.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
  • static.field - 정적 데이터 필드 이름. ! 필수 / ? 선택(기본). Type: string, Default: null, Importance: medium.
  • static.value - 필드 이름이 구성된 경우의 정적 필드 값. Type: string, Default: null, Importance: medium.
  • timestamp.field - 레코드 타임스탬프용 필드 이름. ! 필수 / ? 선택(기본). Type: string, Default: null, Importance: medium.
  • topic.field - Kafka 토픽용 필드 이름. ! 필수 / ? 선택(기본). Type: string, Default: null, Importance: medium.
org.apache.kafka.connect.transforms.InsertHeader

각 레코드에 헤더를 추가합니다.

  • header - 헤더의 이름. Type: string, Valid Values: non-null string, Importance: high.
  • value.literal - 모든 레코드의 헤더 값으로 설정될 리터럴 값. Type: string, Valid Values: non-null string, Importance: high.
org.apache.kafka.connect.transforms.MaskField

지정된 필드를 필드 타입에 맞는 유효한 null 값(0, false, 빈 문자열 등)으로 마스킹합니다.

숫자·문자열 필드에 대해서는 올바른 타입으로 변환되는 선택적 치환 값을 지정할 수 있습니다.

키(MaskField$Key) 또는 값(MaskField$Value)용 구체적 타입을 사용하세요.

  • fields - 마스킹할 필드 이름. Type: list, Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
  • replacement - 모든 'fields' 값(숫자 또는 비어 있지 않은 문자열 값만)에 적용될 커스텀 값 치환. Type: string, Default: null, Valid Values: non-empty string, Importance: low.
org.apache.kafka.connect.transforms.RegexRouter

구성된 정규식과 치환 문자열을 사용해 레코드 토픽을 업데이트합니다.

내부적으로 정규식은 java.util.regex.Pattern으로 컴파일됩니다. 패턴이 입력 토픽과 일치하면 java.util.regex.Matcher#replaceFirst()를 치환 문자열과 함께 사용해 새 토픽을 얻습니다.

  • regex - 일치에 사용할 정규식. Type: string, Valid Values: valid regex, Importance: high.
  • replacement - 치환 문자열. Type: string, Importance: high.
org.apache.kafka.connect.transforms.ReplaceField

필드를 필터링하거나 이름을 변경합니다.

키(ReplaceField$Key) 또는 값(ReplaceField$Value)용 구체적 타입을 사용하세요.

  • exclude - 제외할 필드. 포함할 필드보다 우선합니다. Type: list, Default: "", Importance: medium.
  • include - 포함할 필드. 지정되면 이 필드들만 사용됩니다. Type: list, Default: "", Importance: medium.
  • renames - 필드 이름 변경 매핑. Type: list, Default: "", Valid Values: 콜론으로 구분된 쌍 목록, 예: foo:bar,abc:xyz, Importance: medium.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
org.apache.kafka.connect.transforms.SetSchemaMetadata

레코드의 키(SetSchemaMetadata$Key) 또는 값(SetSchemaMetadata$Value) 스키마에 스키마 이름, 버전 또는 둘 다 설정합니다.

  • schema.name - 설정할 스키마 이름. Type: string, Default: null, Importance: high.
  • schema.version - 설정할 스키마 버전. Type: int, Default: null, Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
org.apache.kafka.connect.transforms.TimestampConverter

Unix epoch, 문자열, Connect Date/Timestamp 타입 같은 서로 다른 형식 간에 타임스탬프를 변환합니다. 개별 필드 또는 전체 값에 적용됩니다.

키(TimestampConverter$Key) 또는 값(TimestampConverter$Value)용 구체적 타입을 사용하세요.

  • target.type - 원하는 타임스탬프 표현: string, unix, Date, Time, Timestamp. Type: string, Valid Values: [string, unix, Date, Time, Timestamp], Importance: high.
  • field - 타임스탬프를 담은 필드, 또는 전체 값이 타임스탬프면 비어 있음. Type: string, Default: "", Importance: high.
  • format - 타임스탬프용 SimpleDateFormat 호환 형식. type=string일 때 출력 생성에, 입력이 문자열일 때 입력 파싱에 사용. Type: string, Default: "", Importance: medium.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.
  • unix.precision - 타임스탬프의 원하는 Unix 정밀도: seconds, milliseconds, microseconds, nanoseconds. type=unix일 때 출력 생성에, 입력이 Long일 때 입력 파싱에 사용. 참고: 이 SMT는 하위 밀리초 성분이 있는 값으로/에서 변환하는 동안 정밀도 손실을 일으킵니다. Type: string, Default: milliseconds, Valid Values: [nanoseconds, microseconds, milliseconds, seconds], Importance: low.
org.apache.kafka.connect.transforms.TimestampRouter

원래 토픽 값과 레코드 타임스탬프의 함수로 레코드의 토픽 필드를 업데이트합니다.

토픽 필드가 대상 시스템의 동일 엔티티 이름(예: 데이터베이스 테이블 또는 검색 인덱스 이름)을 결정하는 데 자주 사용되므로 주로 싱크 커넥터에 유용합니다.

  • timestamp.format - java.text.SimpleDateFormat과 호환되는 타임스탬프 형식 문자열. Type: string, Default: yyyyMMdd, Importance: high.
  • topic.format - 각각 토픽과 타임스탬프의 플레이스홀더로 ${topic}${timestamp}를 포함할 수 있는 형식 문자열. Type: string, Default: ${topic}-${timestamp}, Importance: high.
org.apache.kafka.connect.transforms.ValueToKey

레코드 값의 필드 하위 집합으로 형성된 새 키로 레코드 키를 교체합니다.

  • fields - 레코드 키로 추출할 레코드 값의 필드 이름. Type: list, Importance: high.
  • replace.null.with.default - 기본값이 있는 null 필드를 기본값으로 대체할지 여부. Type: boolean, Default: true, Importance: medium.

프레디케이트 (Predicates)

트랜스포메이션이 어떤 조건을 만족하는 메시지에만 적용되도록 프레디케이트로 구성할 수 있습니다. 특히 Filter 트랜스포메이션과 결합하면 특정 메시지를 선택적으로 걸러내는 데 사용할 수 있습니다.

프레디케이트는 커넥터 구성에 지정됩니다.

  • predicates - 일부 트랜스포메이션에 적용될 프레디케이트의 별칭 집합.
  • predicates.$alias.type - 프레디케이트의 정규화된 클래스 이름.
  • predicates.$alias.$predicateSpecificConfig - 프레디케이트의 구성 속성.

모든 트랜스포메이션에는 predicatenegate라는 암시적 구성 속성이 있습니다. 특정 프레디케이트는 트랜스포메이션의 predicate 구성을 프레디케이트의 별칭으로 설정하면 트랜스포메이션과 연결됩니다. 프레디케이트의 값은 negate 구성 속성을 사용해 반전할 수 있습니다.

예를 들어 많은 다른 토픽에 메시지를 생성하는 소스 커넥터가 있고 다음과 같이 하려 한다고 가정해 봅시다.

  • 'foo' 토픽의 메시지를 완전히 걸러냄.
  • 'bar' 토픽을 제외한 모든 토픽의 레코드에 ExtractField 트랜스포메이션(필드 이름 'other_field') 적용.

이를 위해 먼저 'foo' 토픽으로 향하는 레코드를 걸러내야 합니다. Filter 트랜스포메이션은 이후 처리에서 레코드를 제거하며, TopicNameMatches 프레디케이트를 사용해 특정 정규식과 일치하는 토픽의 레코드에만 트랜스포메이션을 적용할 수 있습니다. TopicNameMatches의 유일한 구성 속성은 토픽 이름과의 일치용 Java 정규식인 pattern입니다. 구성은 다음과 같습니다.

transforms=Filter
transforms.Filter.type=org.apache.kafka.connect.transforms.Filter
transforms.Filter.predicate=IsFoo

predicates=IsFoo
predicates.IsFoo.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsFoo.pattern=foo

다음으로 레코드의 토픽 이름이 'bar'가 아닐 때만 ExtractField를 적용해야 합니다. TopicNameMatches를 직접 사용할 수는 없습니다. 일치하는 토픽 이름이 아닌, 일치하지 않는 토픽 이름에 트랜스포메이션을 적용해야 하기 때문입니다. 트랜스포메이션의 암시적 negate 구성 속성은 프레디케이트가 일치하는 레코드 집합을 반전시키게 해줍니다. 이전 예제에 이 구성의 추가를 반영하면 다음과 같습니다.

transforms=Filter,Extract
transforms.Filter.type=org.apache.kafka.connect.transforms.Filter
transforms.Filter.predicate=IsFoo

transforms.Extract.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.Extract.field=other_field
transforms.Extract.predicate=IsBar
transforms.Extract.negate=true

predicates=IsFoo,IsBar
predicates.IsFoo.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsFoo.pattern=foo

predicates.IsBar.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsBar.pattern=bar

Kafka Connect는 다음 프레디케이트를 포함합니다.

  • TopicNameMatches - 특정 Java 정규식과 일치하는 이름의 토픽에 있는 레코드와 일치.
  • HasHeaderKey - 주어진 키를 가진 헤더가 있는 레코드와 일치.
  • RecordIsTombstone - 톰스톤(tombstone) 레코드, 즉 값이 null인 레코드와 일치.

각 프레디케이트의 구성 방법 세부 사항은 아래에 나열되어 있습니다.

org.apache.kafka.connect.transforms.predicates.HasHeaderKey

구성된 이름의 헤더가 하나 이상 있는 레코드에 대해 참인 프레디케이트.

  • name - 헤더 이름. Type: string, Valid Values: non-empty string, Importance: medium.
org.apache.kafka.connect.transforms.predicates.RecordIsTombstone

톰스톤(값이 null)인 레코드에 대해 참인 프레디케이트.

org.apache.kafka.connect.transforms.predicates.TopicNameMatches

구성된 정규식과 일치하는 토픽 이름을 가진 레코드에 대해 참인 프레디케이트.

  • pattern - 레코드의 토픽 이름과의 일치용 Java 정규식. Type: string, Valid Values: non-empty string, valid regex, Importance: medium.

REST API

Kafka Connect는 서비스로 실행되도록 설계되었으므로 커넥터 관리를 위한 REST API도 제공합니다. 이 REST API는 standalone과 distributed 모드 모두에서 사용할 수 있습니다. REST API 서버는 listeners 구성 옵션으로 구성할 수 있습니다. 이 필드는 protocol://host:port,protocol2://host2:port2 형식의 리스너 목록을 담아야 합니다. 현재 지원되는 프로토콜은 http와 https입니다. 예:

listeners=http://localhost:8080,https://localhost:8443

기본적으로 리스너가 지정되지 않으면 REST 서버는 HTTP 프로토콜로 8083 포트에서 실행됩니다. HTTPS 사용 시 구성에 SSL 구성이 포함되어야 합니다. 기본적으로 ssl.* 설정을 사용합니다. REST API에 Kafka 브로커 연결과 다른 구성을 사용해야 하는 경우 필드에 listeners.https 접두사를 붙일 수 있습니다. 접두사를 사용하면 접두사가 붙은 옵션만 사용되고 접두사 없는 ssl.* 옵션은 무시됩니다. 다음 필드로 REST API의 HTTPS를 구성할 수 있습니다.

  • ssl.keystore.location
  • ssl.keystore.password
  • ssl.keystore.type
  • ssl.key.password
  • ssl.truststore.location
  • ssl.truststore.password
  • ssl.truststore.type
  • ssl.enabled.protocols
  • ssl.provider
  • ssl.protocol
  • ssl.cipher.suites
  • ssl.keymanager.algorithm
  • ssl.secure.random.implementation
  • ssl.trustmanager.algorithm
  • ssl.endpoint.identification.algorithm
  • ssl.client.auth

REST API는 사용자가 Kafka Connect를 모니터링·관리하는 데만 사용되는 것이 아닙니다. distributed 모드에서는 Kafka Connect의 클러스터 간 통신에도 사용됩니다. 팔로워 노드 REST API에서 받은 일부 요청은 리더 노드 REST API로 전달됩니다. 주어진 호스트가 도달 가능한 URI가 리슨하는 URI와 다른 경우 rest.advertised.host.name, rest.advertised.port, rest.advertised.listener 구성 옵션을 사용해 팔로워 노드가 리더와 연결하는 데 사용할 URI를 변경할 수 있습니다. HTTP와 HTTPS 리스너를 모두 사용할 때 rest.advertised.listener 옵션으로 클러스터 간 통신에 사용할 리스너를 정의할 수도 있습니다. 노드 간 통신에 HTTPS를 사용할 때는 동일한 ssl.* 또는 listeners.https 옵션이 HTTPS 클라이언트 구성에 사용됩니다.

현재 지원되는 REST API 엔드포인트는 다음과 같습니다.

  • GET /connectors - 활성 커넥터 목록 반환.
  • POST /connectors - 새 커넥터 생성. 요청 본문은 문자열 name 필드와 커넥터 구성 파라미터가 담긴 객체 config 필드를 포함하는 JSON 객체여야 함. JSON 객체는 선택적으로 STOPPED, PAUSED, RUNNING(기본값) 값을 가질 수 있는 문자열 initial_state 필드도 포함할 수 있음.
  • GET /connectors/{name} - 특정 커넥터에 대한 정보 조회.
  • GET /connectors/{name}/config - 특정 커넥터의 구성 파라미터 조회.
  • PUT /connectors/{name}/config - 특정 커넥터의 구성 파라미터 업데이트.
  • PATCH /connectors/{name}/config - 특정 커넥터의 구성 파라미터 패치. JSON 본문의 null 값은 최종 구성에서 키를 제거함을 나타냄.
  • GET /connectors/{name}/status - 커넥터의 현재 상태 조회. 실행·실패·일시 중지 여부, 할당된 워커, 실패 시 오류 정보, 모든 태스크의 상태 포함.
  • GET /connectors/{name}/tasks - 커넥터에 대해 현재 실행 중인 태스크 목록과 구성 조회.
  • GET /connectors/{name}/tasks/{taskid}/status - 태스크의 현재 상태 조회. 실행·실패·일시 중지 여부, 할당된 워커, 실패 시 오류 정보 포함.
  • PUT /connectors/{name}/pause - 커넥터와 태스크 일시 중지. 커넥터가 재개될 때까지 메시지 처리를 중지. 태스크가 점유한 리소스는 할당된 채로 남으며, 재개 시 커넥터가 빠르게 데이터 처리를 시작할 수 있게 함.
  • PUT /connectors/{name}/stop - 커넥터 정지 및 태스크 종료, 태스크가 점유한 리소스 해제. 리소스 사용 관점에서 일시 중지보다 효율적이지만, 재개 시 데이터 처리 시작이 더 오래 걸릴 수 있음. 커넥터의 오프셋은 정지 상태일 때만 오프셋 관리 엔드포인트로 수정할 수 있음.
  • PUT /connectors/{name}/resume - 일시 중지 또는 정지된 커넥터 재개(일시 중지·정지되지 않았으면 아무것도 안 함).
  • includeTasks 파라미터 - 커넥터 인스턴스와 태스크 인스턴스를 모두 재시작할지(includeTasks=true) 커넥터 인스턴스만 재시작할지(includeTasks=false) 지정. 기본값(false)은 이전 버전과 동일한 동작 유지.
  • onlyFailed 파라미터 - FAILED 상태의 인스턴스만 재시작할지(onlyFailed=true) 모든 인스턴스를 재시작할지(onlyFailed=false) 지정. 기본값(false)은 이전 버전과 동일한 동작 유지.
  • POST /connectors/{name}/tasks/{taskId}/restart - 개별 태스크 재시작(보통 실패했기 때문).
  • DELETE /connectors/{name} - 커넥터 삭제. 모든 태스크 중지 및 구성 삭제.
  • GET /connectors/{name}/topics - 커넥터가 생성된 이후 또는 활성 토픽 집합 리셋 요청 이후 사용 중인 토픽 집합 조회.
  • PUT /connectors/{name}/topics/reset - 커넥터의 활성 토픽 집합을 비우는 요청 전송.
  • GET /connectors/{name}/offsets - 커넥터의 현재 오프셋 조회.
  • DELETE /connectors/{name}/offsets - 커넥터의 오프셋 리셋. 커넥터가 존재하고 정지 상태여야 함(PUT /connectors/{name}/stop 참고).
{
  "offsets": [
    {
      "partition": {
        "filename": "test.txt"
      },
      "offset": {
        "position": 30
      }
    }
  ]
}
  • PATCH /connectors/{name}/offsets - 커넥터의 오프셋 변경. 커넥터가 존재하고 정지 상태여야 함. 요청 본문은 GET /connectors/{name}/offsets 엔드포인트의 응답 본문과 유사한 JSON 배열 offsets 필드를 포함하는 JSON 객체여야 함. FileStreamSourceConnector의 요청 본문 예:
  • PATCH /connectors/{name}/offsets - 커넥터의 오프셋 변경. 커넥터가 존재하고 정지 상태여야 함. 요청 본문은 GET /connectors/{name}/offsets 엔드포인트의 응답 본문과 유사한 JSON 배열 offsets 필드를 포함하는 JSON 객체여야 함. FileStreamSinkConnector의 요청 본문 예:
{
  "offsets": [
    {
      "partition": {
        "kafka_topic": "test",
        "kafka_partition": 0
      },
      "offset": {
        "kafka_offset": 5
      }
    },
    {
      "partition": {
        "kafka_topic": "test",
        "kafka_partition": 1
      },
      "offset": null
    }
  ]
}

특정 파티션의 오프셋을 리셋하려면 "offset" 필드가 null일 수 있습니다(소스·싱크 커넥터 모두에 적용). 소스 커넥터의 경우 요청 본문 형식은 커넥터 구현에 따라 달라지는 반면, 모든 싱크 커넥터에는 공통 형식이 있다는 점에 유의하세요.

Kafka Connect는 커넥터 플러그인에 대한 정보를 얻는 REST API도 제공합니다.

  • GET /connector-plugins - Kafka Connect 클러스터에 설치된 커넥터 플러그인 목록 반환. 이 API는 요청을 처리하는 워커의 커넥터만 검사하므로, 특히 새 커넥터 jar를 추가하는 롤링 업그레이드 중에는 일관되지 않은 결과를 볼 수 있음.
  • GET /connector-plugins/{plugin-type}/config - 지정된 플러그인의 구성 정의 조회.
  • PUT /connector-plugins/{connector-type}/config/validate - 제공된 구성 값을 구성 정의와 대조해 검증. 이 API는 구성별 검증을 수행하고, 검증 중 제안 값과 오류 메시지를 반환.

최상위(root) 엔드포인트에서 지원되는 REST 요청은 다음과 같습니다.

  • GET / - REST 요청을 처리하는 Connect 워커의 버전(소스 코드의 git 커밋 ID 포함)과 연결된 Kafka 클러스터 ID 같은 Kafka Connect 클러스터의 기본 정보 반환.

admin.listeners 구성은 Kafka Connect REST API 서버에서 admin REST API를 구성하는 데 사용할 수 있습니다. listeners 구성과 유사하게 이 필드는 protocol://host:port,protocol2://host2:port2 형식의 리스너 목록을 담아야 합니다. 현재 지원되는 프로토콜은 http와 https입니다. 예:

admin.listeners=http://localhost:8080,https://localhost:8443

기본적으로 admin.listeners가 구성되지 않으면 admin REST API는 일반 리스너에서 사용할 수 있습니다.

현재 지원되는 admin REST API 엔드포인트는 다음과 같습니다.

  • GET /admin/loggers - 수준이 명시적으로 설정된 현재 로거 목록과 로그 수준 조회.
  • GET /admin/loggers/{name} - 지정된 로거의 로그 수준 조회.
  • PUT /admin/loggers/{name} - 지정된 로거의 로그 수준 설정.

admin 로거 REST API에 대한 자세한 내용은 KIP-495를 참고하세요.

Kafka Connect REST API의 전체 규격은 OpenAPI 문서를 참고하세요.

Connect에서의 오류 보고 (Error Reporting in Connect)

Kafka Connect는 처리의 여러 단계에서 발생하는 오류를 처리하기 위한 오류 보고를 제공합니다. 기본적으로 변환 또는 트랜스포메이션 중에 발생한 모든 오류는 커넥터를 실패시킵니다. 각 커넥터 구성은 오류를 건너뛰어 허용할 수도 있으며, 선택적으로 각 오류와 실패한 작업 및 문제 레코드의 세부 정보를(다양한 세부 수준으로) Connect 애플리케이션 로그에 기록할 수도 있습니다. 이러한 메커니즘은 싱크 커넥터가 Kafka 토픽에서 소비한 메시지를 처리할 때의 오류도 잡아내며, 모든 오류는 구성 가능한 "데드 레터 큐"(DLQ) Kafka 토픽에 기록될 수 있습니다.

커넥터의 컨버터, 트랜스포메이션, 또는 싱크 커넥터 자체 내부의 오류를 로그에 보고하려면 커넥터 구성에서 errors.log.enable=true를 설정해 각 오류와 문제 레코드의 토픽·파티션·오프셋 세부 정보를 기록하세요. 추가 디버깅을 위해 errors.log.include.messages=true를 설정해 문제 레코드의 키·값·헤더도 로그에 기록할 수 있습니다(민감 정보를 기록할 수 있음에 유의).

커넥터의 컨버터, 트랜스포메이션, 또는 싱크 커넥터 자체 내부의 오류를 데드 레터 큐 토픽에 보고하려면 errors.deadletterqueue.topic.name을 설정하고, 선택적으로 errors.deadletterqueue.context.headers.enable=true를 설정하세요.

기본적으로 커넥터는 오류나 예외가 발생하면 즉시 "fail fast" 동작을 보입니다. 이는 커넥터 구성에 다음 구성 속성을 기본값과 함께 추가하는 것과 동일합니다.

# disable retries on failure
errors.retry.timeout=0

# do not log the error and their contexts
errors.log.enable=false

# do not record errors in a dead letter queue topic
errors.deadletterqueue.topic.name=

# Fail on first error
errors.tolerance=none

이것들과 다른 관련 커넥터 구성 속성은 다른 동작을 제공하도록 변경할 수 있습니다. 예를 들어 다음 구성 속성을 커넥터 구성에 추가하면 여러 번 재시도하고, 애플리케이션 로그와 my-connector-errors Kafka 토픽에 기록하며, 커넥터 태스크를 실패시키지 않고 보고함으로써 모든 오류를 허용하는 오류 처리를 설정할 수 있습니다.

# retry for at most 10 minutes times waiting up to 30 seconds between consecutive failures
errors.retry.timeout=600000
errors.retry.delay.max.ms=30000

# log error context along with application logs, but do not include configs and messages
errors.log.enable=true
errors.log.include.messages=false

# produce error context into the Kafka topic
errors.deadletterqueue.topic.name=my-connector-errors

# Tolerate all errors.
errors.tolerance=all

정확히 한 번 지원 (Exactly-once support)

Kafka Connect는 싱크 커넥터(0.11.0 버전부터)와 소스 커넥터(3.3.0 버전부터)에 대해 정확히 한 번(exactly-once) 의미론을 제공할 수 있습니다. 정확히 한 번 지원은 실행하는 커넥터의 유형에 크게 의존한다는 점에 유의하세요. 클러스터의 각 노드에 대한 구성에서 모든 올바른 워커 속성을 설정해도, 커넥터가 Kafka Connect 프레임워크의 기능을 활용하도록 설계되지 않았거나 활용할 수 없다면 정확히 한 번이 불가능할 수 있습니다.

싱크 커넥터 (Sink connectors)

싱크 커넥터가 정확히 한 번 의미론을 지원한다면, Connect 워커 수준에서 정확히 한 번을 활성화하려면 그 컨슈머 그룹이 중단된 트랜잭션의 레코드를 무시하도록 구성해야 합니다. 이는 워커 속성 consumer.isolation.levelread_committed로 설정하거나, 이를 지원하는 Kafka Connect 버전을 실행 중이라면 개별 커넥터 구성에서 consumer.override.isolation.level 속성을 read_committed로 설정할 수 있게 하는 커넥터 클라이언트 구성 오버라이드 정책을 사용해 수행할 수 있습니다.

소스 커넥터 (Source connectors)

소스 커넥터가 정확히 한 번 의미론을 지원한다면, 정확히 한 번 소스 커넥터에 대한 프레임워크 수준 지원을 활성화하도록 Connect 클러스터를 구성해야 합니다. 보안 Kafka 클러스터에 대해 실행한다면 추가 ACL이 필요할 수 있습니다. 소스 커넥터에 대한 정확히 한 번 지원은 현재 distributed 모드에서만 사용할 수 있다는 점에 유의하세요. standalone Connect 워커는 정확히 한 번 의미론을 제공할 수 없습니다.

워커 구성 (Worker configuration)

새 Connect 클러스터의 경우 클러스터의 각 노드의 워커 구성에서 exactly.once.source.support 속성을 enabled로 설정하세요. 기존 클러스터의 경우 두 번의 롤링 업그레이드가 필요합니다. 첫 번째 업그레이드 동안에는 exactly.once.source.support 속성을 preparing으로, 두 번째 동안에는 enabled로 설정해야 합니다.

플러그인 발견 (Plugin Discovery)

플러그인 발견은 Connect 워커가 플러그인 클래스를 찾고 이를 커넥터에서 구성·실행할 수 있게 만드는 전략의 이름입니다. 이는 plugin.discovery 워커 구성으로 제어되며, 워커 시작 시간에 큰 영향을 미칩니다. service_load가 가장 빠른 전략이지만, 이 구성을 service_load로 설정하기 전에 플러그인이 호환되는지 확인하는 데 주의해야 합니다.

3.6 버전 이전에는 이 전략을 구성할 수 없었고, 모든 플러그인과 호환되는 only_scan 모드처럼 동작했습니다. 3.6 버전부터 이 모드의 기본값은 hybrid_warn이며, 이것 역시 모든 플러그인과 호환되지만 service_load와 호환되지 않는 플러그인에 대해 경고를 기록합니다. hybrid_fail 전략은 service_load와 호환되지 않는 플러그인이 감지되면 오류로 워커를 중지시켜 모든 플러그인이 호환된다고 주장합니다. 마지막으로 service_load 전략은 다른 모든 모드에서 사용되는 느린 레거시 스캔 메커니즘을 비활성화하고, 대신 더 빠른 ServiceLoader 메커니즘을 사용합니다. 이 메커니즘과 호환되지 않는 플러그인은 사용할 수 없을 수 있습니다.

플러그인 호환성 확인 (Verifying Plugin Compatibility)

모든 플러그인이 service_load와 호환되는지 확인하려면 먼저 Kafka Connect 3.6 이상 버전을 사용 중인지 확인하세요. 그런 다음 다음 확인 중 하나를 수행할 수 있습니다.

  • 기본 hybrid_warn 전략으로 워커를 시작하고, org.apache.kafka.connect 패키지에 대해 WARN 로그를 활성화. plugin.discovery 구성을 언급하는 WARN 로그 메시지가 하나 이상 출력되어야 함. 이 로그 메시지는 모든 플러그인이 호환된다고 명시하거나, 호환되지 않는 플러그인을 나열함.
  • 테스트 환경에서 hybrid_fail로 워커를 시작. 모든 플러그인이 호환되면 시작이 성공. 하나 이상의 플러그인이 호환되지 않으면 워커가 시작에 실패하고, 모든 호환되지 않는 플러그인이 예외에 나열됨.

확인 단계가 성공하면 현재 설치된 플러그인 집합이 호환되는 것이므로 plugin.discovery 구성을 service_load로 변경하는 것이 안전합니다. 확인이 실패하면 service_load 전략을 사용할 수 없으며 호환되지 않는 플러그인 목록을 기록해 두어야 합니다. service_load 전략을 사용하기 전에 모든 플러그인을 처리해야 합니다. 이 확인은 플러그인 설치 또는 버전 변경 후에 수행하는 것이 권장되며, Continuous Integration 환경에서 자동으로 수행할 수 있습니다.

운영자: 아티팩트 마이그레이션 (Operators: Artifact Migration)

Connect 운영자로서 호환되지 않는 플러그인을 발견하면 호환성을 해결할 방법이 여러 가지 있습니다. 선호도가 높은 순서대로 아래에 나열되어 있습니다.

  • 플러그인 제공자의 최신 릴리스를 확인하고, 호환되면 업그레이드.
  • 플러그인 제공자에게 연락해 소스 마이그레이션 지침에 따라 플러그인을 호환되게 마이그레이션하도록 요청한 다음 호환 버전으로 업그레이드.
  • 포함된 마이그레이션 스크립트를 사용해 플러그인 아티팩트를 직접 마이그레이션.

마이그레이션 스크립트는 Kafka 설치의 bin/connect-plugin-path.shbin\windows\connect-plugin-path.bat에 있습니다. 이 스크립트는 JAR 또는 리소스 파일을 추가·수정해 Connect 워커의 plugin.path에 이미 설치된 호환되지 않는 플러그인 아티팩트를 마이그레이션할 수 있습니다. 이는 서명 검증에 실패하도록 아티팩트를 변경할 수 있으므로 코드 서명을 사용하는 환경에는 적합하지 않습니다. 내장 도움말은 --help로 확인하세요.

마이그레이션을 수행하려면 먼저 list 하위 명령으로 스크립트가 사용할 수 있는 플러그인 개요를 얻으세요. 반복 가능한 --worker-config, --plugin-path, --plugin-location 인자로 플러그인을 어디서 찾을지 스크립트에 알려야 합니다. 스크립트는 classpath의 플러그인을 무시하므로, classpath의 커스텀 플러그인은 이 마이그레이션 스크립트와 함께 사용하려면 플러그인 경로로 이동하거나 수동으로 마이그레이션해야 합니다. list 출력과 워커 시작 경고·오류 메시지를 반드시 비교해 영향을 받는 모든 플러그인이 스크립트에 의해 발견되는지 확인하세요.

모든 호환되지 않는 플러그인이 목록에 포함된 것을 확인하면 sync-manifests --dry-run으로 마이그레이션을 드라이 런할 수 있습니다. 이는 마이그레이션 결과를 디스크에 쓰는 것만 제외하고 마이그레이션의 모든 부분을 수행합니다. sync-manifests 명령은 지정된 모든 경로가 쓰기 가능해야 하며 디렉터리 내용을 변경할 수 있다는 점에 유의하세요. 지정된 경로의 플러그인을 백업하거나 쓰기 가능한 디렉터리에 복사하세요.

--dry-run 플래그를 제거하고 실제 마이그레이션을 실행하기 전에 플러그인 백업이 있고 드라이 런이 성공하는지 확인하세요. --dry-run 없이 마이그레이션이 실패하면 부분 마이그레이션된 아티팩트는 버려야 합니다. 마이그레이션은 멱등적(idempotent)이므로 여러 번 실행하거나 이미 마이그레이션된 플러그인에 실행해도 안전합니다. 스크립트가 끝난 후 마이그레이션이 완료되었는지 확인해야 합니다. 마이그레이션 스크립트는 자동 마이그레이션을 위해 Continuous Integration 환경에서 사용하기에 적합합니다.

개발자: 소스 마이그레이션 (Developers: Source Migration)

플러그인을 service_load와 호환되게 만들려면 소스 코드에 ServiceLoader 매니페스트를 추가해야 하며, 이는 릴리스 아티팩트에 패키징되어야 합니다. 매니페스트는 슈퍼클래스 타입의 이름을 딴 META-INF/services/의 리소스 파일이며, 각 줄에 하나씩 정규화된 하위 클래스 이름 목록을 담고 있습니다.

플러그인이 호환되려면 그것이 확장하는 플러그인 슈퍼클래스에 해당하는 매니페스트의 한 줄로 나타나야 합니다. 단일 플러그인이 여러 플러그인 인터페이스를 구현하면 구현하는 각 인터페이스의 매니페스트에 나타나야 합니다. 특정 타입의 플러그인에 대한 클래스가 없다면 그 타입에 대한 매니페스트 파일을 포함할 필요가 없습니다. 플러그인으로 보이지 않아야 하는 클래스는 abstract로 표시해야 합니다. 다음 타입은 매니페스트가 있어야 하는 것으로 예상됩니다.

  • org.apache.kafka.connect.sink.SinkConnector
  • org.apache.kafka.connect.source.SourceConnector
  • org.apache.kafka.connect.storage.Converter
  • org.apache.kafka.connect.storage.HeaderConverter
  • org.apache.kafka.connect.transforms.Transformation
  • org.apache.kafka.connect.transforms.predicates.Predicate
  • org.apache.kafka.common.config.provider.ConfigProvider
  • org.apache.kafka.connect.rest.ConnectRestExtension
  • org.apache.kafka.connect.connector.policy.ConnectorClientConfigOverridePolicy

예를 들어 정규화된 이름이 com.example.MySinkConnector인 커넥터가 하나만 있다면, META-INF/services/org.apache.kafka.connect.sink.SinkConnector의 리소스에 매니페스트 파일 하나만 추가하면 됩니다. 내용은 다음과 유사해야 합니다.

# license header or comment
com.example.MySinkConnector

그런 다음 사전 릴리스(prerelease) 아티팩트로 확인 단계를 사용해 매니페스트가 올바른지 확인해야 합니다. 확인이 성공하면 플러그인을 정상적으로 릴리스할 수 있으며, 운영자는 호환 버전으로 업그레이드할 수 있습니다.

보안 (Security)

Connect에 내재된 보안 문제를 이해하는 것이 중요합니다. 먼저 Connect는 사용자 정의 플러그인 실행을 허용합니다. 이러한 플러그인은 임의 코드를 실행할 수 있으므로 Connect 클러스터에 설치하기 전에 신뢰해야 합니다. 기본적으로 REST API는 보안되지 않으며, 접근할 수 있는 사람이면 누구든 커넥터를 시작·중지할 수 있습니다. REST API를 신뢰할 수 있는 사용자에게만 직접 노출해야 합니다. 그렇지 않으면 Connect 워커에서 임의 코드 실행을 얻기가 쉽습니다. 기본적으로 커넥터는 Connect가 내부적으로 사용하는 Kafka 클라이언트의 구성을 오버라이드할 수도 있습니다. Kafka 4.2.0부터는 connector.client.config.override.policyAllowlist로 설정하는 것이 권장되며, 이는 Kafka 5.0.0부터 기본값이 됩니다. 그리고 필요한 구성만 명시적으로 오버라이드하도록 허용하세요. sasl.jaas.configsasl.login.class 같이 클래스를 로드할 수 있는 구성은, 설계상 Connect 워커에서 코드 실행을 가능하게 하므로 REST API에 신뢰할 수 있는 사용자만 접근할 수 있을 때에만 허용해야 한다는 점을 명심하세요.

ACL 요구 사항

각 Connect 워커의 principal은 다음 ACL이 필요합니다.

작업 리소스 타입 리소스 이름 참고
Read Group Connect 클러스터의 group.id
Read Topic Connect 클러스터의 config.storage.topic
Write Topic Connect 클러스터의 config.storage.topic
Create Topic Connect 클러스터의 config.storage.topic Connect의 config 토픽이 아직 존재하지 않을 때만 필요
Read Topic Connect 클러스터의 offset.storage.topic
Write Topic Connect 클러스터의 offset.storage.topic
Create Topic Connect 클러스터의 offset.storage.topic Connect의 offsets 토픽이 아직 존재하지 않을 때만 필요
Read Topic Connect 클러스터의 status.storage.topic
Write Topic Connect 클러스터의 status.storage.topic
Create Topic Connect 클러스터의 status.storage.topic Connect의 status 토픽이 아직 존재하지 않을 때만 필요
Write TransactionalId connect-cluster-${groupId} (${groupId}는 클러스터의 group.id) 정확히 한 번 지원이 활성화되었거나 exactly.once.source.support가 preparing으로 설정된 경우에만 필요
Describe TransactionalId connect-cluster-${groupId} (${groupId}는 클러스터의 group.id) 정확히 한 번 지원이 활성화되었거나 exactly.once.source.support가 preparing으로 설정된 경우에만 필요
IdempotentWrite Cluster 워커의 config 토픽을 호스팅하는 Kafka 클러스터의 ID 정확히 한 번 지원이 활성화되었거나 exactly.once.source.support가 preparing으로 설정된 경우에만 필요. IdempotentWrite ACL은 2.8부터 폐기되어 pre-2.8 Kafka 클러스터에서 실행되는 Connect 클러스터에만 필요

소스 커넥터를 지원하려면 각 개별 커넥터의 principal에 다음 ACL이 필요합니다.

작업 리소스 타입 리소스 이름 참고
Write Topic 대상으로 사용되는 토픽
Create Topic 대상으로 사용되는 토픽 topic.creation.enable이 true이고 토픽이 아직 존재하지 않을 때만 필요
Describe Topic 커넥터가 사용하는 오프셋 토픽(제공된 경우 커넥터 구성의 offsets.storage.topic 값, 아니면 워커 구성의 offsets.storage.topic 값) 정확히 한 번 지원이 활성화된 경우에만 필요
Write TransactionalId 커넥터가 생성할 각 태스크에 대해 ${groupId}-${connector}-${taskId} (${groupId}는 Connect 클러스터의 group.id, ${connector}는 커넥터 이름, ${taskId}는 0부터 시작하는 태스크 ID) 정확히 한 번 지원이 활성화된 경우에만 필요. 다른 트랜잭션 ID와 충돌 위험이 없거나 충돌이 허용된다면 ${groupId}-${connector}* 와일드카드 접두사를 편의상 사용 가능
Describe TransactionalId 커넥터가 생성할 각 태스크에 대해 ${groupId}-${connector}-${taskId} 정확히 한 번 지원이 활성화된 경우에만 필요. 와일드카드 접두사 사용 가능
IdempotentWrite Cluster 워커의 config 토픽을 호스팅하는 Kafka 클러스터의 ID 정확히 한 번 지원이 활성화된 경우에만 필요. IdempotentWrite ACL은 2.8부터 폐기되어 pre-2.8 Kafka 클러스터에만 필요

싱크 커넥터를 지원하려면 각 개별 커넥터의 principal에 다음 ACL이 필요합니다.

작업 리소스 타입 리소스 이름 참고
Read Group connect-${connector} (${connector}는 커넥터 이름, 또는 Connect 구성에 consumer.group.id가 있으면 그 값, 또는 커넥터 구성에 consumer.overrides.group.id가 있으면 그 값)
Read Topic 커넥터가 소비할 싱크 토픽 커넥터의 topics 또는 topics.regex 옵션으로 식별
Write Topic 커넥터의 errors.deadletterqueue.topic.name errors.deadletterqueue.topic.name이 비어 있지 않은 값으로 설정된 경우에만 필요

일부 커넥터는 Kafka Connect 프레임워크가 관리하지 않는 Kafka 토픽을 추가로 사용합니다(예: change data capture 커넥터는 스키마 히스토리 저장에 추가 토픽을 사용). 필요한 추가 ACL의 세부 사항은 커넥터 문서를 참조하세요. 이러한 ACL은 개별 커넥터의 principal에 추가할 수 있습니다.

더 알아보기 (Learn more)

  • 트랜스포메이션과 프레디케이트 조합은 데이터 다듬기와 라우팅을 강력하게 만들어요.
  • 정확히 한 번 지원은 커넥터 유형과 클러스터 구성에 따라 달라지니 잘 확인해야 해요.
  • Connect 보안과 ACL 요구 사항을 꼼꼼히 지키면 운영 리스크를 줄일 수 있어요.