다음 표에는 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 데이터의 경우 |
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
이 구성은 다음을 수행합니다.
org.apache.kafka.connect.transforms.predicates.HasHeaderKey유형의hasKey라는 술어를 구성합니다. 이 조건자는 키가DoNotProcess인 헤더가 포함된 모든 메시지와 일치합니다.org.apache.kafka.connect.transforms.Filter유형의dropMessage이라는 변환을 구성합니다. 이 변환은 구성된 조건자와 일치하는 모든 메시지를 삭제합니다.변환을 술어
hasKey에 연결합니다. 이렇게 하면DoNotProcess헤더 키가 있는 메시지만 변환에 의해 삭제됩니다.