Kafka 커넥터 문제 해결
Kafka 커넥터 문제 해결 (Troubleshooting the Kafka connector)
중요:
- 사전 공지: 클래식 Kafka 커넥터(v3 이하)는 현재 완전히 지원되지만, 향후 폐기될 예정이에요.
- 조치: 즉시 변경할 필요는 없어요. 현재 워크로드는 안전하며 계속 완전히 지원돼요.
- 일정: Snowflake는 2026년 중반에 공식 폐기 공지를 발표할 계획이에요. 공지 후 수명 종료까지 18개월의 마이그레이션 기간이 시작돼요.
- 권장사항: 모든 새 구현에는 Snowflake Connector for Kafka (v4)를 사용하세요.
마이그레이션 지침은 v3에서 v4로 마이그레이션을 참고하세요.
이 섹션은 Kafka 커넥터로 데이터를 수집할 때 발생하는 문제를 해결하는 방법을 설명해요.
출처: 문서
본문
오류 알림
Snowpipe에 대한 오류 알림을 구성하세요. Snowpipe가 로드 중 파일 오류를 만나면 이 기능은 구성된 클라우드 메시징 서비스로 알림을 푸시해 데이터 파일을 분석할 수 있게 해줘요. 자세한 내용은 Snowpipe 오류 알림을 참고하세요.
일반 문제 해결 단계
Kafka 커넥터를 사용한 로드 문제를 해결하려면 다음 단계를 완료하세요.
1단계: 테이블의 COPY 기록 보기
대상 테이블의 로드 활동 기록을 쿼리하세요. 자세한 내용은 COPY_HISTORY 뷰를 참고하세요. COPY_HISTORY 출력에 예상된 파일 세트가 포함되지 않으면 더 이른 기간을 쿼리하세요. 파일이 이전 파일의 복제본이었다면 로드 기록이 원본 파일 로드 시도 당시의 활동을 기록했을 수 있어요. STATUS 컬럼은 특정 파일 세트가 로드되었는지, 부분적으로 로드되었는지, 로드에 실패했는지를 나타내요. FIRST_ERROR_MESSAGE 컬럼은 시도가 부분적으로 로드되었거나 실패했을 때 이유를 제공해요.
Kafka 커넥터는 로드할 수 없는 파일을 대상 테이블과 연결된 스테이지로 이동시켜요. 테이블 스테이지를 참조하는 구문은 @[namespace.]%table_name이에요.
LIST를 사용해 테이블 스테이지에 있는 모든 파일을 나열하세요.
예를 들어:
LIST @mydb.public.%mytable;
파일 이름은 다음 형식 중 하나예요. 각 형식을 만드는 조건은 표에 설명돼 있어요:
| 파일 유형 | 설명 |
| 원시 바이트 (Raw bytes) | 이 파일들은 다음 패턴과 일치해요: <connector_name>/<table_name>/<partition>/offset_(<key>/<value>_)<timestamp>.gz. 이 파일들의 경우 Kafka 레코드를 원시 바이트에서 소스 파일 형식(Avro, JSON, Protobuf)으로 변환할 수 없었어요. 이 문제의 일반적인 원인은 레코드에서 문자가 손실된 네트워크 실패예요. Kafka 커넥터가 더 이상 원시 바이트를 파싱할 수 없어 깨진 레코드가 됐어요. |
| 소스 파일 형식 (Avro, JSON, Protobuf) | 이 파일들은 다음 패턴과 일치해요: <connector_name>/<table_name>/<partition>/<start_offset>_<end_offset>_<timestamp>.<file_type>.gz. 이 파일들의 경우 Kafka 커넥터가 원시 바이트를 소스 파일 형식으로 변환한 후 Snowpipe가 오류를 만나 파일을 로드할 수 없었어요. |
다음 섹션은 각 파일 유형의 문제를 해결하는 방법을 설명해요.
원시 바이트 (Raw bytes)
파일 이름 <connector_name>/<table_name>/<partition>/offset_(<key>/<value>_)<timestamp>.gz는 원시 바이트에서 소스 파일 형식으로 변환되지 않은 레코드의 정확한 오프셋을 포함해요. 문제를 해결하려면 레코드를 새 레코드로 Kafka 커넥터에 다시 보내세요.
소스 파일 형식 (Avro, JSON, Protobuf)
Snowpipe가 Kafka 토픽용으로 만든 내부 스테이지의 파일에서 데이터를 로드할 수 없다면 Kafka 커넥터는 파일을 대상 테이블의 스테이지로 소스 파일 형식으로 이동시켜요.
파일 세트에 여러 문제가 있으면 COPY_HISTORY 출력의 FIRST_ERROR_MESSAGE 컬럼은 첫 번째로 만난 오류만 나타내요. 파일의 모든 오류를 보려면 테이블 스테이지에서 파일을 검색하고, 이를 명명된 스테이지에 업로드한 다음, copy 옵션 VALIDATION_MODE를 RETURN_ALL_ERRORS로 설정한 COPY INTO <table> 문을 실행해야 해요. VALIDATION_MODE copy 옵션은 COPY 문이 로드할 데이터를 검증하고 지정된 검증 옵션에 따라 결과를 반환하도록 지시해요. 이 copy 옵션을 지정하면 데이터는 로드되지 않아요. 문에서 Kafka 커넥터로 로드하려고 시도했던 파일 세트를 참조하세요.
데이터 파일의 문제를 해결하면 하나 이상의 COPY 문으로 데이터를 수동으로 로드할 수 있어요.
다음 예제는 mydb.public 데이터베이스와 스키마의 mytable 테이블에 대한 테이블 스테이지에 있는 데이터 파일을 참조해요.
테이블 스테이지의 데이터 파일을 검증하고 오류를 해결하려면:
LIST를 사용해 테이블 스테이지에 있는 모든 파일을 나열하세요.- 예:
LIST @mydb.public.%mytable;
- 예:
- 이 섹션의 예제는 데이터 파일의 소스 형식이 JSON이라고 가정해요.
GET을 사용해 Kafka 커넥터가 만든 파일을 로컬 머신에 다운로드하세요.- 예를 들어 파일을 로컬 머신의
data라는 디렉토리에 다운로드:- Linux 또는 macOS:
GET @mydb.public.%mytable file:///data/; - Microsoft Windows:
GET @mydb.public.%mytable file://C:\data\;
- Linux 또는 macOS:
- 예를 들어 파일을 로컬 머신의
- 소스 Kafka 파일과 같은 형식의 데이터 파일을 저장하는 명명된 내부 스테이지를
CREATE STAGE로 만들기.- 예를 들어 JSON 파일을 저장하는
kafka_json이라는 내부 스테이지 생성:CREATE STAGE kafka_json FILE_FORMAT = (TYPE = JSON);
- 예를 들어 JSON 파일을 저장하는
PUT을 사용해 테이블 스테이지에서 다운로드한 파일을 업로드하세요.- 예를 들어 로컬 머신의
data디렉토리에 다운로드한 파일 업로드:- Linux 또는 macOS:
PUT file:///data/ @mydb.public.kafka_json; - Microsoft Windows:
PUT file://C:\data\ @mydb.public.kafka_json;
- Linux 또는 macOS:
- 예를 들어 로컬 머신의
- 테스트 목적의 두 variant 컬럼이 있는 임시 테이블을 만드세요. 테이블은 스테이징된 데이터 파일을 검증하는 데만 사용돼요. 테이블에는 데이터가 로드되지 않아요. 테이블은 현재 사용자 세션이 끝나면 자동으로 삭제돼요:
CREATE TEMPORARY TABLE t1 (col1 variant); COPY INTO table … VALIDATION_MODE = 'RETURN_ALL_ERRORS'문을 실행해 데이터 파일에서 만난 모든 오류를 검색하세요. 이 문은 지정된 스테이지의 파일을 검증해요. 테이블에는 데이터가 로드되지 않아요:
COPY INTO mydb.public.t1
FROM @mydb.public.kafka_json
FILE_FORMAT = (TYPE = JSON)
VALIDATION_MODE = 'RETURN_ALL_ERRORS';
- 로컬 머신의 데이터 파일에서 보고된 모든 오류를 수정하세요.
PUT을 사용해 수정된 파일을 테이블 스테이지 또는 명명된 내부 스테이지에 업로드하세요.- 다음 예제는 파일을 테이블 스테이지에 업로드하고 기존 파일을 덮어써요:
- Linux 또는 macOS:
PUT file:///tmp/myfile.csv @mydb.public.%mytable OVERWRITE = TRUE; - Windows:
PUT file://C:\temp\myfile.csv @mydb.public.%mytable OVERWRITE = TRUE;
- Linux 또는 macOS:
- 다음 예제는 파일을 테이블 스테이지에 업로드하고 기존 파일을 덮어써요:
VALIDATION_MODE옵션 없이COPY INTO table을 사용해 데이터를 대상 테이블에 로드하세요. 데이터가 성공적으로 로드되면 스테이지에서 데이터 파일을 삭제하는 copy 옵션PURGE = TRUE를 선택적으로 사용하거나,REMOVE로 테이블 스테이지에서 파일을 수동으로 삭제할 수 있어요:
COPY INTO mydb.public.mytable(RECORD_METADATA, RECORD_CONTENT)
FROM (SELECT $1:meta, $1:content FROM @mydb.public.%mytable)
FILE_FORMAT = (TYPE = 'JSON')
PURGE = TRUE;
2단계: Kafka 커넥터 로그 파일 분석
COPY_HISTORY 뷰에 데이터 로드 기록이 없으면 Kafka 커넥터의 로그 파일을 분석하세요. 커넥터는 이벤트를 로그 파일에 써요. Snowflake Kafka 커넥터는 모든 Kafka 커넥터 플러그인과 같은 로그 파일을 공유한다는 점에 유의하세요. 이 로그 파일의 이름과 위치는 Kafka Connect 구성 파일에 있어야 해요. 자세한 내용은 Apache Kafka 소프트웨어에 제공된 문서를 참고하세요.
Kafka 커넥터 로그 파일에서 Snowflake 관련 오류 메시지를 검색하세요. 대부분의 메시지에는 ERROR 문자열이 있고 파일 이름 com.snowflake.kafka.connector...을 포함해 이 메시지를 더 쉽게 찾을 수 있어요.
만날 수 있는 가능한 오류는:
-
구성 오류 (Configuration error): 오류의 가능한 원인:
- 커넥터가 토픽을 구독할 적절한 정보를 가지고 있지 않음.
- 커넥터가 Snowflake 테이블에 쓸 적절한 정보를 가지고 있지 않음 (예: 인증용 키 페어가 잘못됨).
Kafka 커넥터는 파라미터를 검증한다는 점에 유의하세요. 커넥터는 호환되지 않는 각 구성 파라미터에 대해 오류를 던져요. 오류 메시지는 Kafka Connect 클러스터의 로그 파일에 기록돼요. 구성 문제를 의심하면 해당 로그 파일의 오류를 확인하세요.
-
읽기 오류 (Read error): 커넥터가 다음 이유로 Kafka에서 읽지 못했을 수 있어요:
- Kafka 또는 Kafka Connect가 실행 중이 아닐 수 있음.
- 메시지가 아직 전송되지 않았을 수 있음.
- 메시지가 삭제(만료)되었을 수 있음.
-
쓰기 오류 — 스테이지 (Write error - stage): 오류의 가능한 원인:
- 스테이지에 대한 권한 부족.
- 스테이지 공간 부족.
- 스테이지가 삭제됨.
- 다른 사용자나 프로세스가 스테이지에 예상치 못한 파일을 썼음.
-
쓰기 오류 — 테이블 (Write error - table): 오류의 가능한 원인:
- 테이블에 대한 권한 부족.
3단계: Kafka Connect 확인
Kafka Connect 로그 파일에 오류가 보고되지 않으면 Kafka Connect를 확인하세요. 문제 해결 지침은 Apache Kafka 소프트웨어 공급업체가 제공한 문서를 참고하세요.
특정 문제 해결
같은 토픽 파티션과 오프셋을 가진 중복 행
버전 1.4(이상)의 Kafka 커넥터로 데이터를 로드할 때 대상 테이블의 같은 토픽 파티션·오프셋을 가진 중복 행은 로드 작업이 기본 실행 시간 제한인 300000밀리초(300초)를 초과했음을 나타낼 수 있어요. 원인을 확인하려면 Kafka Connect 로그 파일에서 다음 오류를 확인하세요:
org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member.
This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records.
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.sendOffsetCommitRequest(ConsumerCoordinator.java:1061)
오류를 해결하려면 Kafka 구성 파일(예: <kafka_dir>/config/connect-distributed.properties)에서 다음 속성 중 하나를 변경하세요:
consumer.max.poll.interval.ms— 실행 시간 제한을900000(900초)으로 늘림.consumer.max.poll.records— 각 작업으로 로드되는 레코드 수를50으로 줄임.
스트리밍 채널 오프셋 마이그레이션 응답 실패 오류 코드: 5023
v2.1.0(이상) 커넥터 버전으로 업그레이드할 때 Snowpipe Streaming 채널 이름 형식에 변경이 도입됐어요. 결과적으로 이전에 커밋된 오프셋에 대한 정보를 감지하는 로직이 이전에 커밋된 정보를 찾지 못해요. 이는 다음과 같은 예외로 나타나요:
com.snowflake.kafka.connector.internal.SnowflakeKafkaConnectorException: [SF_KAFKA_CONNECTOR] Exception: Failure in Streaming Channel Offset Migration Response Error Code: 5023
Detail: Streaming Channel Offset Migration from Source to Destination Channel has no/invalid response, please contact Snowflake Support
Message: Snowflake experienced a transient exception, please retry the migration request.
이 오류를 해결하려면 Kafka 구성 파일(예: <kafka_dir>/config/connect-distributed.properties)에 다음 구성 속성을 추가하세요:
enable.streaming.channel.offset.migration—false로 설정해 자동 오프셋 마이그레이션을 비활성화.
여러 토픽을 지원하도록 커넥터 구성
단일 Kafka 커넥터 인스턴스가 각각 여러 파티션을 가진 많은 수의 토픽을 지원하는 문제를 만난 적이 있어요. 커넥터의 구성이 유효해 보였음에도 Snowflake로 아무 데이터도 수집할 수 없는 끝없는 리밸런스 주기가 발생했어요. 이 문제는 Snowpipe Streaming 수집 모드(snowflake.ingestion.method=SNOWPIPE_STREAMING)에 특정했지만, 지침은 Snowpipe 수집 모드(snowflake.ingestion.method=SNOWPIPE)에도 적용돼요.
이 문제는 로그 파일에 다음 로그 메시지가 반복적으로 기록되는 것으로 나타나요:
[Worker-xyz] [timestamp] INFO [my-connector|task-id] [SF_INGEST] Channel is marked as closed
이것은 보통 커넥터가 regex로 토픽을 수집하도록 구성했을 때 발생할 수 있어요. Kafka 구성 파일(예: <kafka_dir>/config/connect-distributed.properties)에 다음 옵션 세트를 적용하는 것을 권장해요:
consumer.override.partition.assignment.strategy— 파티션 할당 전략을 태스크에 대해org.apache.kafka.clients.consumer.CooperativeStickyAssignor로 구성. 이렇게 하면 수집된 채널이 사용 가능한 태스크에 고르게 분배되어 리밸런싱 위험이 줄어들어요.CooperativeStickyAssignor는 이 알려진 이슈 때문에 Kafka Connect 버전 3.0.1 이상이 필요하다는 점에 유의하세요.tasks.max— 커넥터당 인스턴스화된 태스크 수는 사용 가능한 CPU 수를 초과하지 않아야 해요. 기본 드라이버는 사용 가능한 CPU를 기반으로 스로틀링 메커니즘을 구현해요. 동시 요청 수를 늘리면 시스템의 메모리 압력이 증가하지만 삽입 처리 시간도 길어져 커넥터의 하트비트 누락으로 직접 이어질 수 있어요.
커넥터의 타임아웃 값에 대해 직접 영향을 주는 구성 속성 세트가 있어요:
consumer.override.heartbeat.interval.ms— 모니터 스레드(각 태스크에 하나씩 연결됨)가 Kafka에 하트비트를 보내는 빈도를 정의해요. 기본값은3000ms지만, 시스템 부하가 높은 경우5000ms로 늘려 실험해 볼 수 있어요.consumer.override.session.timeout.ms— 브로커가 소비자를 무효 상태로 간주하고 리밸런스를 시도하기 전에 기다리는 시간을 정의해요. 이 설정은 보통 하트비트 간격의 3배여야 해요. 하트비트를5000ms로 구성했다면 이 값을15000ms로 설정하세요.consumer.override.max.poll.interval.ms— 기본 Kafka의poll()호출 사이의 최대 간격을 정의해요. 폴 사이의 시간은 기본적으로 커넥터가 데이터 배치를 처리하는 시간(Snowflake 업로드와 커밋 포함)에 해당해요. 여러 태스크가 데이터를 처리하는 시나리오에서는 기본 Snowflake 연결이 요청을 스로틀링하기 시작해 처리 시간이 길어질 수 있어요. 시나리오에 따라, 특히 수집할 큰 초기 레코드 수로 커넥터를 시작할 때 이 값을 20분(1200000ms)까지 늘릴 수 있어요.consumer.override.rebalance.timeout.ms— 리밸런스가 발생할 때 태스크당 채널 수가 많은 시나리오에서는 처리 재개 위치를 파악하기 위해 채널당 많은 기본 로직이 있어요. 이 코드는 순차적으로 실행되므로 태스크당 채널 수가 많을수록 초기 설정이 오래 걸려요. 각 채널이 초기화를 완료할 수 있도록 이 속성을 충분히 큰 값으로 구성하세요. 3분(180000ms) 값이 좋은 시작점이에요.
커넥터에 사용 가능한 힙 메모리를 인지하는 것도 중요해요. 여러 커넥터가 동시에 실행되거나 하나의 커넥터가 여러 토픽에서 데이터를 수집하는 시나리오에서 특히 중요해요. 각 토픽의 파티션은 단일 채널에 매핑되므로 메모리가 필요해요.
Xmx 설정을 통해 Kafka Connect 프로세스 메모리 설정을 조정하세요. 한 가지 방법은 KAFKA_OPTS 환경 변수를 정의하고 그에 맞게 설정하는 것이에요 (즉, KAFKA_OPTS=-Xmx4G).
파일 클리너가 예기치 않게 파일 제거
SNOWPIPE와 함께 Kafka 커넥터를 사용할 때 여러 토픽에서 단일 테이블로 데이터를 수집하는 문제를 만날 수 있어요. 구성에 snowflake.topic2table.map 항목이 없거나 토픽과 테이블 사이에 1:1 매핑이 있으면 이 문제는 적용되지 않아요.
Kafka 커넥터는 스테이지에 업로드할 레코드가 있는 파일을 생성해요. 이 파일은 다음 패턴에 따라 형식이 지정돼요:
snowflake_kafka_connector_<connector-name>_stage_<table-name>/<connector-name>/<table-name>/<partition-id>/<low-watermark>_<high-watermark>_<timestamp>.json.gz
문제는 <partition-id>에 있어요. 여러 토픽이 단일 테이블로 데이터를 로드하면 partition-id 값에 중복이 있을 가능성이 높아요. 이는 정상적인 커넥터 운영에서는 문제가 아니에요. 하지만 커넥터가 재시작되거나 리밸런스되면 클리너 프로세스가 스테이지에 로드된(아직 수집되지 않은) 파일을 잘못된 파티션과 잘못 연결해 삭제하기로 결정할 수 있으며, 이는 데이터 손실 이벤트로 이어질 수 있어요.
버전 2.5.0의 커넥터는 partition-id에 소스 토픽의 해시코드를 포함해 단일 토픽 파티션과 정확히 일치하는 고유한 파일 이름을 보장함으로써 이 문제를 해결해요. 이 수정은 기본적으로 활성화돼 있으며(snowflake.snowpipe.stageFileNameExtensionEnabled), snowflake.topic2table.map에서 대상 테이블이 두 번 이상 나열된 구성에만 영향을 줘요.
구성이 이 기능의 영향을 받으면 스테이지에 오래된 파일이 업로드될 수 있어요. 커넥터가 시작될 때 스테이지에 그러한 파일이 있는지 확인해요. NOTE: For table로 시작하고 감지된 파일 목록이 이어지는 로그 항목을 찾아야 해요.
스테이지에서 영향을 받는 파일이 있는지 수동으로 확인할 수도 있어요:
- 영향을 받는 스테이지 찾기:
show stages like 'snowflake_kafka_connector%<your table name>';
- 스테이지 파일 나열:
list @<your stage name> pattern = '.+/<your-table-name>/[0-9]{1,4}/[0-9]+_[0-9]+_[0-9]+\.json\.gz$';
위 명령은 테이블의 스테이지와 일치하고 파티션 ID가 0-9999 범위인 모든 파일을 나열해요. 이 파일들은 더 이상 수집되지 않으므로 다운로드하거나 삭제할 수 있어요.
문제 보고
Snowflake Support에 도움을 요청할 때 다음 파일을 준비하세요:
- Kafka 커넥터용 구성 파일. (중요: 파일을 Snowflake에 제공하기 전에 프라이빗 키를 제거하세요.)
- Kafka 커넥터 로그 사본. 파일에 기밀 또는 민감한 정보가 포함되지 않도록 하세요.
- JDBC 로그 파일. 로그 파일을 생성하려면 Kafka 커넥터를 실행하기 전에 Kafka Connect 클러스터에
JDBC_TRACE = true환경 변수를 설정하세요. JDBC 로그 파일에 대한 자세한 내용은 Snowflake Community의 이 문서를 참고하세요. - 연결(Connect) 로그 파일. 로그 파일을 생성하려면
etc/kafka/connect-log4j.properties파일을 편집하고log4j.appender.stdout.layout.ConversionPattern속성을 다음과 같이 설정하세요:
log4j.appender.stdout.layout.ConversionPattern=[%d] %p %X{connector.context}%m (%c:%L)%n
커넥터 컨텍스트는 Kafka 버전 2.3 이상에서 사용할 수 있어요. 자세한 내용은 Confluent 웹사이트의 로깅 개선 정보를 참고하세요.