Apache Kafka 커넥터

Apache Kafka 커넥터

Flink는 정확히-한 번(exactly-once) 보장으로 Kafka 토픽에서 데이터를 읽고 쓰기 위한 Apache Kafka 커넥터를 제공해요.

출처: Apache Kafka Connector

본문

Dependency

Apache Flink는 최신 버전의 Kafka 클라이언트를 추적하는 범용(universal) Kafka 커넥터와 함께 제공돼요. 사용하는 클라이언트의 버전은 Flink 릴리스마다 달라질 수 있어요. 최신 Kafka 클라이언트는 브로커 버전 2.1.0 이상과 역호환돼요. Kafka 호환성에 대한 자세한 내용은 공식 Kafka 문서를 참고하세요.

Flink 버전 2.3용 커넥터는 아직 제공되지 않아요.

Flink의 스트리밍 커넥터는 바이너리 배포판의 일부가 아니에요. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.

PyFlink 잡에서 사용하려면 다음 의존성이 필요해요:

flink-connector-kafka Flink 버전 2.3용 SQL jar는 아직 제공되지 않아요.

Kafka Source

이 부분은 새로운 데이터 소스 API 기반의 Kafka source를 설명해요.

Usage

Kafka source는 KafkaSource 인스턴스를 구성하기 위한 빌더 클래스를 제공해요. 아래 코드 스니펫은 "input-topic" 토픽의 earliest offset에서, consumer group "my-group"으로, 메시지의 value만 string으로 역직렬화하는 KafkaSource를 빌드하는 방법을 보여줘요.

KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers(brokers)
    .setTopics("input-topic")
    .setGroupId("my-group")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");
source = KafkaSource.builder() \
    .set_bootstrap_servers(brokers) \
    .set_topics("input-topic") \
    .set_group_id("my-group") \
    .set_starting_offsets(KafkaOffsetsInitializer.earliest()) \
    .set_value_only_deserializer(SimpleStringSchema()) \
    .build()

env.from_source(source, WatermarkStrategy.no_watermarks(), "Kafka Source")

KafkaSource를 빌드하는 데 필요한 속성은 다음과 같아요:

  • Bootstrap servers, setBootstrapServers(String)으로 구성
  • 구독할 Topics / partitions — 아래 Topic-partition subscription 참고
  • Kafka 메시지를 파싱할 Deserializer — 아래 Deserializer 참고

Topic-partition Subscription

Kafka source는 3가지 토픽-파티션 구독 방식을 제공해요:

  • Topic list: 토픽 목록의 모든 파티션에서 메시지를 구독. 예: Java KafkaSource.builder().setTopics("topic-a", "topic-b"); Python KafkaSource.builder().set_topics("topic-a", "topic-b")
  • Topic pattern: 제공된 정규식과 이름이 일치하는 모든 토픽에서 메시지를 구독. 예: Java KafkaSource.builder().setTopicPattern("topic.*"); Python KafkaSource.builder().set_topic_pattern("topic.*")
  • Partition set: 제공된 파티션 집합의 파티션을 구독. 예: Java final HashSet<TopicPartition> partitionSet = new HashSet<>(Arrays.asList(new TopicPartition("topic-a", 0), new TopicPartition("topic-b", 5))); KafkaSource.builder().setPartitions(partitionSet); Python partition_set = {KafkaTopicPartition("topic-a", 0), KafkaTopicPartition("topic-b", 5)} KafkaSource.builder().set_partitions(partition_set)

Deserializer

Kafka 메시지를 파싱하려면 역직렬화기가 필요해요. 역직렬화기(Deserialization schema)는 setDeserializer(KafkaRecordDeserializationSchema)로 구성할 수 있으며, KafkaRecordDeserializationSchema는 Kafka ConsumerRecord를 역직렬화하는 방법을 정의해요.

Kafka ConsumerRecord의 value만 필요하면 빌더에서 setValueOnlyDeserializer(DeserializationSchema)를 사용할 수 있고, DeserializationSchema는 Kafka 메시지 value의 바이너리를 역직렬화하는 방법을 정의해요.

Kafka 메시지 value 역직렬화에 Kafka Deserializer도 사용할 수 있어요. 예를 들어 StringDeserializer로 value를 string으로 역직렬화:

import org.apache.kafka.common.serialization.StringDeserializer;

KafkaSource.<String>builder()
        .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class));

현재 PyFlink는 Kafka 레코드 value의 역직렬화를 커스터마이즈하기 위해 set_value_only_deserializer만 지원해요.

KafkaSource.builder().set_value_only_deserializer(SimpleStringSchema())

Starting Offset

Kafka source는 OffsetsInitializer를 지정해 서로 다른 offset에서 메시지 소비를 시작할 수 있어요. 내장 initializer는:

KafkaSource.builder()
    // Start from committed offset of the consuming group, without reset strategy
    .setStartingOffsets(OffsetsInitializer.committedOffsets())
    // Start from committed offset, also use EARLIEST as reset strategy if committed offset doesn't exist
    .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
    // Start from the first record whose timestamp is greater than or equals a timestamp (milliseconds)
    .setStartingOffsets(OffsetsInitializer.timestamp(1657256176000L))
    // Start from earliest offset
    .setStartingOffsets(OffsetsInitializer.earliest())
    // Start from latest offset
    .setStartingOffsets(OffsetsInitializer.latest());
KafkaSource.builder() \
    # Start from committed offset of the consuming group, without reset strategy
    .set_starting_offsets(KafkaOffsetsInitializer.committed_offsets()) \
    # Start from committed offset, also use EARLIEST as reset strategy if committed offset doesn't exist
    .set_starting_offsets(KafkaOffsetsInitializer.committed_offsets(KafkaOffsetResetStrategy.EARLIEST)) \
    # Start from the first record whose timestamp is greater than or equals a timestamp (milliseconds)
    .set_starting_offsets(KafkaOffsetsInitializer.timestamp(1657256176000)) \
    # Start from the earliest offset
    .set_starting_offsets(KafkaOffsetsInitializer.earliest()) \
    # Start from the latest offset
    .set_starting_offsets(KafkaOffsetsInitializer.latest())

내장 initializer가 요구사항을 충족하지 못하면 커스텀 offsets initializer를 구현할 수도 있어요. (PyFlink에서는 미지원)

offsets initializer를 지정하지 않으면 기본적으로 OffsetsInitializer.earliest()가 사용돼요.

Boundedness

Kafka source는 스트리밍과 배치 실행 모드 모두를 지원하도록 설계됐어요. 기본적으로 KafkaSource는 스트리밍 방식으로 실행되도록 설정되어 Flink 잡이 실패하거나 취소될 때까지 절대 멈추지 않아요. setBounded(OffsetsInitializer)를 사용해 stopping offset을 지정하고 source를 배치 모드로 실행할 수 있어요. 모든 파티션이 stopping offset에 도달하면 source가 종료돼요.

setUnbounded(OffsetsInitializer)를 사용해 KafkaSource를 스트리밍 모드로 실행하면서 stopping offset에서 멈추도록 할 수도 있어요. 모든 파티션이 지정된 stopping offset에 도달하면 source가 종료돼요.

Additional Properties

위에서 설명한 속성 외에도 setProperties(Properties)setProperty(String, String)KafkaSourceKafkaConsumer에 임의의 속성을 설정할 수 있어요. KafkaSource는 구성에 대한 다음 옵션이 있어요:

  • client.id.prefix — Kafka consumer의 client ID에 사용할 접두사 정의
  • partition.discovery.interval.ms — Kafka source가 새 파티션을 발견하는 간격(밀리초) 정의. 아래 Dynamic Partition Discovery 참고
  • register.consumer.metrics — Flink metric group에 KafkaConsumer의 메트릭을 등록할지 지정
  • commit.offsets.on.checkpoint — 체크포인트에 소비 offset을 Kafka 브로커에 커밋할지 지정

KafkaConsumer 구성에 대해서는 Apache Kafka 문서를 참고하세요.

구성되어 있어도 빌더가 다음 키를 덮어쓴다는 점에 유의하세요:

  • auto.offset.reset.strategy — starting offsets에 대해 OffsetsInitializer#getAutoOffsetResetStrategy()로 덮어씀
  • partition.discovery.interval.mssetBounded(OffsetsInitializer)가 호출되면 -1로 덮어씀

Dynamic Partition Discovery

토픽 확장 또는 토픽 생성 같은 시나리오를 Flink 잡을 재시작하지 않고 처리하기 위해, Kafka source는 제공된 토픽-파티션 구독 패턴 아래에서 새 파티션을 주기적으로 발견하도록 구성할 수 있어요. 파티션 발견을 활성화하려면 partition.discovery.interval.ms 속성에 양수 값을 설정하세요:

KafkaSource.builder()
    .setProperty("partition.discovery.interval.ms", "10000"); // discover new partitions per 10 seconds
KafkaSource.builder() \
    .set_property("partition.discovery.interval.ms", "10000")  # discover new partitions per 10 seconds

파티션 발견 간격은 기본적으로 5분이에요. 이 기능을 비활성화하려면 파티션 발견 간격을 명시적으로 양수가 아닌 값으로 설정해야 해요.

Event Time and Watermarks

기본적으로 레코드는 Kafka ConsumerRecord에 내장된 타임스탬프를 이벤트 타임으로 사용해요. 레코드 자체에서 이벤트 타임을 추출하고 워터마크를 하류로 내보내도록 자체 WatermarkStrategy를 정의할 수 있어요:

env.fromSource(kafkaSource, new CustomWatermarkStrategy(), "Kafka Source With Custom Watermark Strategy");

WatermarkStrategy를 정의하는 방법은 이 문서에서 설명해요. (PyFlink 미지원)

Idleness

병렬도가 파티션 수보다 높으면 Kafka Source는 자동으로 idle 상태가 되지 않아요. 병렬도를 낮추거나 워터마크 전략에 idle timeout을 추가해야 해요. 특정 시간 동안 스트림의 한 파티션에 레코드가 흐르지 않으면 그 파티션은 "idle"로 간주되고 하류 연산자의 워터마크 진행을 막지 않아요.

WatermarkStrategy#withIdleness를 정의하는 방법은 이 문서에서 설명해요.

Consumer Offset Committing

Kafka source는 Flink의 체크포인트 상태와 Kafka 브로커에 커밋된 offset 사이의 일관성을 보장하기 위해 체크포인트가 완료될 때 현재 소비 offset을 커밋해요.

체크포인팅이 활성화되지 않으면 Kafka source는 Kafka consumer의 내부 자동 주기 offset 커밋 로직에 의존하며, 이는 Kafka consumer 속성의 enable.auto.commitauto.commit.interval.ms로 구성돼요.

Kafka source는 내결함성에 커밋된 offset에 의존하지 않는다는 점에 유의하세요. offset을 커밋하는 것은 모니터링을 위해 consumer와 consuming group의 진행 상황을 노출하기 위한 것뿐이에요.

Monitoring

Kafka source는 해당 스코프에서 다음 메트릭을 노출해요.

Scope of Metric

Scope Metrics User Variables Description Type
Operator currentEmitEventTimeLag n/a 레코드 이벤트 타임스탬프부터 source 커넥터¹가 레코드를 내보낸 시간까지의 시간 간격: currentEmitEventTimeLag = EmitTime - EventTime. Gauge
watermarkLag n/a 워터마크가 벽시계 시간보다 뒤처지는 시간 간격: watermarkLag = CurrentTime - Watermark Gauge
sourceIdleTime n/a source가 어떤 레코드도 처리하지 않은 시간 간격: sourceIdleTime = CurrentTime - LastRecordProcessTime Gauge
pendingRecords n/a source가 아직 가져오지 않은 레코드 수. 예: Kafka 파티션에서 consumer offset 이후의 사용 가능한 레코드. Gauge
KafkaSourceReader.commitsSucceeded n/a offset 커밋이 켜져 있고 체크포인팅이 활성화된 경우 Kafka에 대한 성공적인 offset 커밋 총 수. Counter
KafkaSourceReader.commitsFailed n/a offset 커밋이 켜져 있고 체크포인팅이 활성화된 경우 Kafka에 대한 offset 커밋 실패 총 수. Kafka에 offset을 다시 커밋하는 것은 consumer 진행 상황을 노출하는 수단일 뿐이므로 커밋 실패는 Flink의 체크포인트 파티션 offset 무결성에 영향을 주지 않아요. Counter
KafkaSourceReader.committedOffsets topic, partition 파티션별로 Kafka에 마지막으로 성공적으로 커밋된 offset. 특정 파티션 메트릭은 토픽 이름과 파티션 id로 지정할 수 있어요. Gauge
KafkaSourceReader.currentOffsets topic, partition 파티션별로 consumer의 현재 읽기 offset. 특정 파티션 메트릭은 토픽 이름과 파티션 id로 지정할 수 있어요. Gauge

¹ 이 메트릭은 마지막으로 처리된 레코드에 대해 기록된 순간 값이에요. 지연 히스토그램은 비쌀 수 있으므로 이 메트릭이 제공돼요. 순간 지연 값은 보통 지연의 좋은 지표가 돼요.

Kafka Consumer Metrics

Kafka consumer의 모든 메트릭도 KafkaSourceReader.KafkaConsumer 그룹 아래에 등록돼요. 예를 들어 Kafka consumer 메트릭 "records-consumed-total"은 <some_parent_groups>.operator.KafkaSourceReader.KafkaConsumer.records-consumed-total 메트릭으로 보고돼요.

register.consumer.metrics 옵션을 구성해 Kafka consumer의 메트릭 등록 여부를 구성할 수 있어요. 이 옵션은 기본적으로 true로 설정돼요.

Kafka consumer의 메트릭에 대해서는 Apache Kafka 문서를 참고하세요.

javax.management.InstanceAlreadyExistsException: kafka.consumer:[...]을 포함하는 스택 트레이스가 있는 경고가 나타나면, 같은 client.id로 여러 KafkaConsumer를 등록하려 하고 있을 가능성이 있어요. 이 경고는 모든 사용 가능한 메트릭이 메트릭 시스템에 올바르게 전달되지 않았음을 나타내요. 모든 KafkaSource에 다른 client.id.prefix를 구성하고 잡의 다른 KafkaConsumer가 같은 client.id를 사용하지 않도록 해야 해요.

Security

암호화와 인증을 포함한 보안 구성을 활성화하려면 Kafka source에 추가 속성으로 보안 구성을 설정하면 돼요. 아래 코드 스니펫은 Kafka source를 SASL 메커니즘으로 PLAIN을 사용하고 JAAS 구성을 제공하도록 구성하는 방법을 보여줘요:

KafkaSource.builder()
    .setProperty("security.protocol", "SASL_PLAINTEXT")
    .setProperty("sasl.mechanism", "PLAIN")
    .setProperty("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"username\" password=\"password\";");
KafkaSource.builder() \
    .set_property("security.protocol", "SASL_PLAINTEXT") \
    .set_property("sasl.mechanism", "PLAIN") \
    .set_property("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"username\" password=\"password\";")

더 복잡한 예시로, 보안 프로토콜로 SASL_SSL을 사용하고 SASL 메커니즘으로 SCRAM-SHA-256을 사용:

KafkaSource.builder()
    .setProperty("security.protocol", "SASL_SSL")
    // SSL configurations
    // Configure the path of truststore (CA) provided by the server
    .setProperty("ssl.truststore.location", "/path/to/kafka.client.truststore.jks")
    .setProperty("ssl.truststore.password", "test1234")
    // Configure the path of keystore (private key) if client authentication is required
    .setProperty("ssl.keystore.location", "/path/to/kafka.client.keystore.jks")
    .setProperty("ssl.keystore.password", "test1234")
    // SASL configurations
    // Set SASL mechanism as SCRAM-SHA-256
    .setProperty("sasl.mechanism", "SCRAM-SHA-256")
    // Set JAAS configurations
    .setProperty("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";");
KafkaSource.builder() \
    .set_property("security.protocol", "SASL_SSL") \
    # SSL configurations
    # Configure the path of truststore (CA) provided by the server
    .set_property("ssl.truststore.location", "/path/to/kafka.client.truststore.jks") \
    .set_property("ssl.truststore.password", "test1234") \
    # Configure the path of keystore (private key) if client authentication is required
    .set_property("ssl.keystore.location", "/path/to/kafka.client.keystore.jks") \
    .set_property("ssl.keystore.password", "test1234") \
    # SASL configurations
    # Set SASL mechanism as SCRAM-SHA-256
    .set_property("sasl.mechanism", "SCRAM-SHA-256") \
    # Set JAAS configurations
    .set_property("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";")

잡 JAR에서 Kafka 클라이언트 의존성을 리로케이트하면 sasl.jaas.config의 login module 클래스 경로가 다를 수 있으므로 JAR 안의 실제 모듈 클래스 경로로 다시 써야 할 수도 있어요.

보안 구성에 대한 자세한 설명은 Apache Kafka 문서의 "Security" 섹션을 참고하세요.

Kafka Rack Awareness

Kafka rack awareness를 사용하면 Rack ID에 기반해 Kafka consumer가 읽을 클라우드 지역과 가용 영역을 Flink가 선택·제어할 수 있어요. 이 기능은 consumer가 가장 가까운 Kafka 브로커(같은 클라우드 지역·가용 영역에 함께 배치될 가능성이 높음)에 연결할 수 있게 하므로 네트워크 비용과 지연을 줄여줘요. 클라이언트의 rack은 client.rack 구성으로 표시되며, 브로커의 broker.rack 구성과 일치해야 해요.

https://kafka.apache.org/documentation/#consumerconfigs_client.rack

RackId

setRackIdSupplier()는 consumer의 rack을 결정할 수 있게 해주는 Builder 메서드예요. 제공되면 Supplier는 Task Manager에서 consumer가 설정될 때 실행되고, consumer의 client.rack 구성이 그 값으로 설정돼요.

이를 구현하는 방법 중 하나는 taskManager 안의 환경 변수에 setRackId를 같게 하는 것이에요:

.setRackIdSupplier(() -> System.getenv("TM_NODE_AZ"))

"TM_NODE_AZ"는 사용하려는 zone을 포함하는 TaskManager 컨테이너의 환경 변수 이름이에요.

Behind the Scene

Kafka source가 새 데이터 소스 API의 설계에서 어떻게 작동하는지 관심이 있다면 이 부분을 참고로 읽을 수 있어요. 새 데이터 소스 API에 대한 자세한 내용은 데이터 소스 문서와 FLIP-27이 더 설명적인 논의를 제공해요.

새 데이터 소스 API 추상화 아래에서 Kafka source는 다음 컴포넌트로 구성돼요:

Source Split

Kafka source의 source split은 Kafka 토픽의 파티션을 나타내요. Kafka source split은 다음으로 구성돼요:

  • 나타내는 TopicPartition
  • 파티션의 시작 offset (Starting offset)
  • 파티션의 중지 offset (Stopping offset) — source가 bounded mode로 실행될 때만 사용 가능

Kafka source split의 상태는 파티션의 현재 소비 offset도 저장하며, Kafka source reader가 snapshot될 때 상태가 불변 split으로 변환되어 현재 offset이 불변 split의 시작 offset으로 할당돼요.

자세한 내용은 KafkaPartitionSplitKafkaPartitionSplitState 클래스를 확인하세요.

Split Enumerator

Kafka의 split enumerator는 제공된 토픽 파티션 구독 패턴 아래에서 새 split(파티션)을 발견하고, round-robin 방식으로 subtask 전체에 고르게 분포시켜 split을 reader에 할당하는 책임이 있어요. Kafka source의 split enumerator는 source reader에게 split을 적극적으로(eagerly) 밀어주므로 source reader의 split 요청을 처리할 필요가 없어요.

Source Reader

Kafka source의 source reader는 제공된 SourceReaderBase를 확장하며, 하나의 SplitReader가 구동하는 하나의 KafkaConsumer로 여러 할당된 split(파티션)을 읽는 single-thread-multiplexed 스레드 모델을 사용해요. 메시지는 Kafka에서 가져온 직후 SplitReader에서 역직렬화돼요. split의 상태, 즉 메시지 소비의 현재 진행 상황은 KafkaRecordEmitter에 의해 갱신되며, 레코드가 하류로 내보내질 때 이벤트 타임을 할당하는 책임도 있어요.

Kafka SourceFunction

FlinkKafkaConsumer는 deprecated이며 Flink 1.17에서 제거될 예정이에요. KafkaSource를 사용하세요.

옛 참조는 Flink 1.13 documentation에서 볼 수 있어요.

Kafka Sink

KafkaSink는 하나 이상의 Kafka 토픽에 레코드 스트림을 기록할 수 있게 해줘요.

Usage

Kafka sink는 KafkaSink 인스턴스를 구성하기 위한 빌더 클래스를 제공해요. 아래 코드 스니펫은 at-least-once 전달 보장으로 String 레코드를 Kafka 토픽에 기록하는 방법을 보여줘요.

DataStream<String> stream = ...;
        
KafkaSink<String> sink = KafkaSink.<String>builder()
        .setBootstrapServers(brokers)
        .setRecordSerializer(KafkaRecordSerializationSchema.builder()
            .setTopic("topic-name")
            .setValueSerializationSchema(new SimpleStringSchema())
            .build()
        )
        .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
        .build();
        
stream.sinkTo(sink);
sink = KafkaSink.builder() \
    .set_bootstrap_servers(brokers) \
    .set_record_serializer(
        KafkaRecordSerializationSchema.builder()
            .set_topic("topic-name")
            .set_value_serialization_schema(SimpleStringSchema())
            .build()
    ) \
    .set_delivery_guarantee(DeliveryGuarantee.AT_LEAST_ONCE) \
    .build()

stream.sink_to(sink)

KafkaSink을 빌드하는 데 필요한 속성은 다음과 같아요:

  • Bootstrap servers, setBootstrapServers(String)
  • Record serializer, setRecordSerializer(KafkaRecordSerializationSchema)
  • DeliveryGuarantee.EXACTLY_ONCE로 delivery guarantee를 구성하면 setTransactionalIdPrefix(String)도 사용해야 해요

Serializer

들어오는 요소를 데이터 스트림에서 Kafka producer 레코드로 변환하려면 항상 KafkaRecordSerializationSchema를 제공해야 해요. Flink는 key/value 직렬화, 토픽 선택, 파티셔닝 같은 일반적인 빌딩 블록을 제공하는 스키마 빌더를 제공해요. 더 많은 제어를 위해 인터페이스를 직접 구현할 수도 있어요.

KafkaRecordSerializationSchema.builder()
    .setTopicSelector((element) -> {<your-topic-selection-logic>})
    .setValueSerializationSchema(new SimpleStringSchema())
    .setKeySerializationSchema(new SimpleStringSchema())
    .setPartitioner(new FlinkFixedPartitioner())
    .build();
KafkaRecordSerializationSchema.builder() \
    .set_topic_selector(lambda element: <your-topic-selection-logic>) \
    .set_value_serialization_schema(SimpleStringSchema()) \
    .set_key_serialization_schema(SimpleStringSchema()) \
    # set partitioner is not supported in PyFlink
    .build()

항상 value 직렬화 메서드와 토픽(선택 메서드)을 설정해야 해요. 또한 setKafkaKeySerializer(Serializer) 또는 setKafkaValueSerializer(Serializer)로 Flink serializer 대신 Kafka serializer를 사용할 수도 있어요.

Fault Tolerance

전반적으로 KafkaSink는 세 가지 다른 DeliveryGuarantee를 지원해요. DeliveryGuarantee.AT_LEAST_ONCEDeliveryGuarantee.EXACTLY_ONCE의 경우 Flink의 체크포팅이 활성화되어야 해요. 기본적으로 KafkaSinkDeliveryGuarantee.NONE을 사용해요. 아래에서 각 보장에 대한 설명을 볼 수 있어요.

  • DeliveryGuarantee.NONE은 어떤 보장도 제공하지 않아요: Kafka 브로커에 문제가 있을 때 메시지가 손실될 수 있고, Flink 오류가 발생하면 메시지가 중복될 수 있어요.
  • DeliveryGuarantee.AT_LEAST_ONCE: sink는 체크포인트에서 Kafka 버퍼의 모든 대기 중인 레코드가 Kafka producer에 의해 승인될 때까지 기다려요. Kafka 브로커에 문제가 있어도 메시지가 손실되지 않지만, Flink가 오래된 입력 레코드를 다시 처리하므로 Flink가 재시작될 때 메시지가 중복될 수 있어요.
  • DeliveryGuarantee.EXACTLY_ONCE: 이 모드에서 KafkaSink는 체크포인트에서 Kafka에 커밋될 Kafka 트랜잭션에 모든 메시지를 기록해요. 따라서 consumer가 커밋된 데이터만 읽으면(Kafka consumer 구성 isolation.level 참고) Flink 재시작 시 중복이 보이지 않아요. 그러나 이는 체크포인트가 기록될 때까지 레코드 가시성을 지연시키므로 체크포인트 기간을 그에 맞게 조정해야 해요. 같은 Kafka 클러스터에서 실행되는 애플리케이션 전체에 고유한 transactionalIdPrefix를 사용해 여러 실행 중인 잡이 트랜잭션을 방해하지 않도록 해야 해요! 또한 Kafka 트랜잭션 타임아웃(카프카 producer transaction.timeout.ms)을 최대 체크포인트 기간 + 최대 재시작 기간보다 크게 조정하는 것을 강력히 권장해요. 그렇지 않으면 Kafka가 커밋되지 않은 트랜잭션을 만료시킬 때 데이터 손실이 발생할 수 있어요.

Monitoring

Kafka sink는 해당 스코프에서 다음 메트릭을 노출해요.

Scope Metrics User Variables Description Type
Operator currentSendTime n/a 마지막 레코드를 보내는 데 걸린 시간. 이 메트릭은 마지막으로 처리된 레코드에 대해 기록된 순간 값이에요. Gauge

Kafka Producer

FlinkKafkaProducer는 deprecated이며 Flink 1.15에서 제거될 예정이에요. KafkaSink를 사용하세요.

옛 참조는 Flink 1.13 documentation에서 볼 수 있어요.

Kafka Connector Metrics

Flink의 Kafka 커넥터는 커넥터의 동작을 분석하기 위해 Flink의 메트릭 시스템을 통해 일부 메트릭을 제공해요. producer와 consumer는 모든 지원 버전에 대해 Kafka의 내부 메트릭을 Flink의 메트릭 시스템을 통해 내보내요. Kafka 문서가 내보내진 모든 메트릭을 나열해요.

Kafka 메트릭 전달을 비활성화하는 것도 가능한데, KafkaSource의 경우 이 섹션에서 설명한 register.consumer.metrics를 구성하거나, KafkaSink를 사용할 때 producer 속성을 통해 register.producer.metrics를 false로 설정하면 돼요.

Enabling Kerberos Authentication

Flink는 Kafka 커넥터를 통해 Kerberos로 구성된 Kafka 설치에 인증하기 위한 일급(first-class) 지원을 제공해요. 간단히 flink-conf.yaml에서 Flink를 구성해 Kafka에 대한 Kerberos 인증을 활성화하면 돼요:

  • Kerberos 자격 증명 구성 —
    • security.kerberos.login.use-ticket-cache: 기본적으로 true이며, Flink는 kinit이 관리하는 티켓 캐시의 Kerberos 자격 증명을 사용하려 시도해요. YARN에 배포된 Flink 잡에서 Kafka 커넥터를 사용할 때 티켓 캐시를 사용한 Kerberos 인증은 작동하지 않는다는 점에 유의하세요.
    • security.kerberos.login.keytabsecurity.kerberos.login.principal: 대신 Kerberos keytab을 사용하려면 이 두 속성 모두에 값을 설정하세요.
  • KafkaClientsecurity.kerberos.login.contexts에 추가: 이는 Flink가 구성된 Kerberos 자격 증명을 Kafka 인증에 사용할 Kafka login context에 제공하도록 지시해요.

Kerberos 기반 Flink 보안이 활성화되면, 내부 Kafka 클라이언트에 전달되는 제공된 속성 구성에 다음 두 설정을 포함하기만 하면 Flink Kafka Consumer 또는 Producer로 Kafka에 인증할 수 있어요:

  • security.protocolSASL_PLAINTEXT(기본값 NONE)로 설정: Kafka 브로커와 통신하는 데 사용되는 프로토콜. 독립형 Flink 배포를 사용할 때는 SASL_SSL도 사용할 수 있어요. SSL용 Kafka 클라이언트 구성 방법은 여기를 참고하세요.
  • sasl.kerberos.service.namekafka(기본값 kafka)로 설정: 이 값은 Kafka 브로커 구성에 사용된 sasl.kerberos.service.name과 일치해야 해요. 클라이언트와 서버 구성 사이의 service name 불일치는 인증 실패를 일으켜요.

Flink의 Kerberos 보안 구성에 대한 자세한 내용은 여기를 참고하세요. Flink가 내부적으로 Kerberos 기반 보안을 설정하는 방법에 대한 자세한 내용은 여기에서 찾을 수 있어요.

Upgrading to the Latest Connector Version

일반적인 업그레이드 단계는 잡과 Flink 버전 업그레이드 가이드에 설명되어 있어요. Kafka의 경우 다음 단계도 따라야 해요:

  • Flink와 Kafka 커넥터 버전을 동시에 업그레이드하지 마세요.
  • Consumer에 group.id가 구성되어 있는지 확인하세요.
  • read offset이 Kafka에 커밋되도록 consumer에서 setCommitOffsetsOnCheckpoints(true)를 설정하세요. 중지하고 savepoint를 만들기 전에 이 작업을 하는 것이 중요해요. 이 설정을 활성화하려면 옛 커넥터 버전에서 stop/restart 주기를 해야 할 수도 있어요.
  • Kafka에서 read offset을 얻도록 consumer에서 setStartFromGroupOffsets(true)를 설정하세요. 이는 Flink 상태에 read offset이 없을 때만 적용되며, 그래서 다음 단계가 매우 중요해요.
  • source/sink의 할당된 uid를 변경하세요. 이렇게 하면 새 source/sink가 옛 source/sink 연산자에서 상태를 읽지 않게 보장돼요.
  • savepoint에 여전히 이전 커넥터 버전의 상태가 있으므로 --allow-non-restored-state로 새 잡을 시작하세요.

Troubleshooting

Flink를 사용할 때 Kafka에 문제가 있으면, Flink는 KafkaConsumerKafkaProducer를 감쌀 뿐이며 문제가 Flink와 무관할 수 있고 Kafka 브로커 업그레이드, Kafka 브로커 재구성, 또는 Flink에서 KafkaConsumer/KafkaProducer 재구성으로 해결될 수 있음을 기억하세요. 흔한 문제의 몇 가지 예가 아래에 나열돼 있어요.

Data loss

Kafka 구성에 따라 Kafka가 기록을 승인한 후에도 데이터 손실이 발생할 수 있어요. 특히 Kafka 구성에서 다음 속성을 기억하세요:

  • acks
  • log.flush.interval.messages
  • log.flush.interval.ms
  • log.flush.*

위 옵션들의 기본값은 쉽게 데이터 손실로 이어질 수 있어요. 자세한 설명은 Kafka 문서를 참조하세요.

UnknownTopicOrPartitionException

이 오류의 한 가능한 원인은 새 leader 선거가 진행 중일 때입니다. 예를 들어 Kafka 브로커를 재시작한 후나 재시작 중에 발생해요. 이는 재시도 가능한 예외이므로 Flink 잡은 재시작하고 정상 동작을 재개할 수 있어야 해요. producer 설정에서 retries 속성을 변경해 피할 수도 있어요. 그러나 이는 메시지 순서 변경을 일으킬 수 있고, 원치 않으면 max.in.flight.requests.per.connection을 1로 설정해 피할 수 있어요.

ProducerFencedException

이 예외의 원인은 브로커 측의 트랜잭션 타임아웃일 가능성이 가장 높아요. KAFKA-6119의 구현으로 (producerId, epoch)는 트랜잭션 타임아웃 후 펜싱(fenced)되고 모든 대기 중인 트랜잭션이 중단돼요(각 transactional.id는 단일 producerId에 매핑됨. 자세한 내용은 다음 블로그 포스트 참고).

더 알아보기 (Learn more)