커넥터 개요

다음 표에는 Managed Service for Apache Kafka에서 지원하는 Kafka Connect 커넥터 유형이 나와 있습니다. 이러한 커넥터를 사용하여 Apache Kafka를 애플리케이션 및 기타 Google Cloud 서비스와 통합할 수 있습니다.

커넥터 설명 사용 사례
MirrorMaker 2.0 한 Kafka 클러스터에서 다른 Kafka 클러스터로 주제 및 데이터를 복제합니다. 데이터 복제, 재해 복구, 데이터 마이그레이션
BigQuery 싱크 Kafka 주제에서 BigQuery 테이블로 데이터를 스트리밍합니다. 데이터 웨어하우징, 분석
Cloud Storage 싱크 Kafka 주제에서 Cloud Storage 버킷으로 데이터를 스트리밍합니다. 데이터 레이크 수집, 데이터 보관처리
Pub/Sub 싱크 Kafka 주제에서 Pub/Sub 주제로 데이터를 스트리밍합니다. 서비스 통합, 실시간 알림
Pub/Sub 소스 Pub/Sub 구독에서 Kafka 주제로 메시지를 스트리밍합니다. 실시간 데이터 수집, 이벤트 기반 아키텍처

전환수

변환기는 Kafka 레코드 데이터의 직렬화 및 역직렬화를 담당합니다. Kafka 주제에 있는 원시 바이트 형식과 Kafka Connect에서 사용하는 내부 구조화된 데이터 표현 간에 변환합니다.

  • 싱크 커넥터의 경우 변환기는 주제의 와이어 형식에서 Kafka Connect 내부 데이터 형식으로 데이터를 역직렬화합니다. 커넥터는 이 형식을 사용하여 타겟 시스템에 씁니다.

  • 소스 커넥터의 경우 변환기는 Kafka Connect 내부 데이터 형식의 데이터를 Kafka 주제의 지정된 와이어 형식으로 직렬화합니다.

변환기는 커넥터가 외부 시스템과 호환되는 형식으로 Kafka 레코드를 읽거나 쓰도록 합니다.

커넥터를 구성할 때 다음 속성을 설정합니다.

  • 키 변환기 (key.converter): Kafka 레코드 키를 직렬화하고 역직렬화할 때 사용할 변환기입니다.

  • 값 변환기 (value.converter): Kafka 레코드 값을 직렬화하고 역직렬화할 때 사용할 변환기입니다.

변환기를 지정하지 않으면 기본 변환기 유형은 org.apache.kafka.connect.converters.ByteArrayConverter이며, 이는 원시 바이트 형식으로 데이터를 전달합니다.

지원되는 변환기

Managed Service for Apache Kafka는 다음 기본 제공 변환기를 지원합니다.

전환 액션을 취한 사용자형식
io.confluent.connect.avro.AvroConverter Apache Avro
org.apache.kafka.connect.converters.BooleanConverter 불리언
org.apache.kafka.connect.converters.ByteArrayConverter

바이트 배열

기본 변환기 유형입니다. 두 시스템에서 메시지의 정확한 콘텐츠를 유지합니다.

org.apache.kafka.connect.converters.DoubleConverter Double
org.apache.kafka.connect.converters.FloatConverter 부동 소수점 수
org.apache.kafka.connect.converters.IntegerConverter 정수
org.apache.kafka.connect.json.JsonConverter

JSON

스키마가 없는 JSON 데이터의 경우 value.converter.schemas.enable=false도 설정합니다.

org.apache.kafka.connect.converters.LongConverter 길이
org.apache.kafka.connect.converters.ShortConverter 짧은 버전
org.apache.kafka.connect.storage.StringConverter 문자열

변환기 선택은 커넥터 유형과 Kafka에 저장하는 데이터에 따라 달라집니다. 자세한 내용은 특정 커넥터의 문서를 참고하세요.

작업

커넥터는 병렬로 작동하는 하나 이상의 작업을 만들어 데이터를 전송합니다. 커넥터가 생성하는 태스크 수의 상한을 설정하려면 커넥터의 tasks.max 구성 속성을 설정하세요. 커넥터가 이 값보다 적은 수의 작업을 만들 수 있습니다.

tasks.max 값을 늘리면 처리량이 향상되지만 리소스 소비 (CPU 및 메모리)도 증가합니다. 최적의 값은 워크로드와 Connect 클러스터 작업자에 할당된 리소스에 따라 다릅니다. 싱크 커넥터의 경우 Kafka 주제 파티션 수도 병렬 처리에 영향을 줄 수 있습니다.

작업 재시작 정책

실패가 발생할 때의 동작을 결정하는 커넥터의 작업 재시작 정책을 설정할 수 있습니다. 커넥터는 다음 정책을 지원합니다.

  • 다시 시작하지 않음 커넥터는 실패한 작업을 다시 시작하지 않습니다. 이 정책은 기본 동작입니다. 디버깅이나 오류 후 수동 개입이 필요한 상황에 유용합니다.

  • 지수 백오프로 다시 시작합니다. 커넥터는 지연 (백오프 기간이라고 함) 후 실패한 작업을 다시 시작합니다. 지연 시간은 후속 실패마다 기하급수적으로 증가합니다. 이 정책은 대부분의 프로덕션 워크로드에 권장됩니다.

    지수 백오프 정책을 사용하는 경우 최소 및 최대 백오프 값도 설정하세요. 최소 백오프는 60초보다 커야 하고 최대 백오프는 7, 200초보다 작아야 합니다.

변환 및 술어

Managed Service for Apache Kafka는 기본 Kafka Connect 변환프레디케이트를 지원합니다.

변환을 사용하면 Managed Service for Apache Kafka (소스 커넥터의 경우) 또는 외부 시스템 (싱크 커넥터의 경우)으로 전송되기 전에 개별 메시지를 수정할 수 있습니다. 변환을 사용하여 민감한 정보를 마스킹하거나, 타임스탬프를 추가하거나, 필드 이름을 바꿀 수 있습니다.

프레디케이트를 사용하면 특정 조건에 따라 데이터를 필터링하여 메시지 속성에 따라 변환이 적용되는 메시지를 결정할 수 있습니다.

예를 들어 DoNotProcess 헤더 키가 포함된 메시지를 무시하도록 싱크 커넥터를 구성하려면 다음 구성을 추가합니다.

transforms=dropMessage
transforms.dropMessage.type=org.apache.kafka.connect.transforms.Filter
transforms.dropMessage.predicate=hasKey
predicates=hasKey
predicates.hasKey.type=org.apache.kafka.connect.transforms.predicates.HasHeaderKey
predicates.hasKey.name=DoNotProcess

이 구성은 다음을 수행합니다.

  1. org.apache.kafka.connect.transforms.predicates.HasHeaderKey 유형의 hasKey라는 술어를 구성합니다. 이 조건자는 키가 DoNotProcess인 헤더가 포함된 모든 메시지와 일치합니다.

  2. org.apache.kafka.connect.transforms.Filter 유형의 dropMessage이라는 변환을 구성합니다. 이 변환은 구성된 조건자와 일치하는 모든 메시지를 삭제합니다.

  3. 변환을 술어 hasKey에 연결합니다. 이렇게 하면 DoNotProcess 헤더 키가 있는 메시지만 변환에 의해 삭제됩니다.

다음 단계

Apache Kafka®는 미국 및/또는 다른 국가에서 사용되는 Apache Software Foundation 또는 해당 계열사의 등록 상표입니다.