본문으로 건너뛰기

Apache Kafka

외부 데이터 streaming platform (Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API) 의 topic 을 게이트웨이가 클라이언트로 outbound 접속 해 consumer subscribe / producer send 합니다. 도착한 record 는 in-memory 캐시에 저장되고, 태그의 수집 주기마다 캐시값을 읽어 갑니다.

상황어떤 모드를 쓰나요
외부 IT 시스템이 REST POST 로 값을 보내올 때HTTP 푸시
외부 서버가 WebSocket 으로 push 할 때WebSocket Client
외부 IIoT MQTT broker 가 topic 으로 push 할 때MQTT Client
Kafka topic 을 consumer 로 받을 때Apache Kafka (이 페이지)
사내 PLC 의 값을 직접 읽을 때Modbus / OPC-UA

등록 폼 입력값

입력 칸무엇을 적나요예시
IP 주소Kafka broker 호스트 (bootstrap)kafka.example.com, 10.0.0.50
포트broker 포트 (평문 9092, TLS 9093)9092, 9093
GROUP IDconsumer group id (옵션)plant-floor-A, edge-line1
SECURITY보안 프로토콜PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL
SASL MECHSASL 메커니즘 (옵션)PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512
SASL JAASSASL 인증 JAAS 설정 (옵션)org.apache.kafka.common.security.plain.PlainLoginModule required username="u" password="p";
AUTO OFFSET신규 group 시작점latest (기본) / earliest
수집 주기캐시된 값을 게이트웨이가 읽어가는 주기 (ms)1000

실제 broker URL: <host>:<port> (예: kafka.example.com:9092). bootstrap 서버는 한 호스트만 등록해도 broker 가 cluster 정보를 알아서 알려줍니다.

GROUP ID 미지정 시 plantpulse-edge-<opc_id> 형식으로 자동 생성됩니다.


태그의 PLC 주소 표기 — 4-mode JSON

태그의 PLC 주소 = Kafka topic + 4-mode 디코더. MQTT / WebSocket 과 동일한 spec 입니다.

모드형식동작
SCALARsensors/temp 또는 sensors/temp.value메시지 전체를 String 으로. 메시지가 JSON object 면 raw 폴백.
KEYsensors/temp:temperaturetop-level JSON key 의 값 (예: {"temperature":25.3,"humidity":60}25.3)
PATHsensors/temp:$.data.tags.T1JSON Pointer 동적 evaluate (중첩 key 지원)
RAWsensors/temp:_raw_마지막 record value 전체 (디버깅)

read 의 첫 호출은 lazy subscribe (빈 문자열 반환). 다음 polling cycle 부터 캐시값이 들어옵니다.


자주 쓰는 사례

사례어떻게
사내 Kafka clusterhost = bootstrap broker IP, 9092 (평문)
Confluent Cloudhost = pkc-XXX.region.aws.confluent.cloud, 9092 + SASL_SSL + PLAIN mechanism + API key/secret JAAS
AWS MSK (IAM auth)host = b-1.<cluster>...amazonaws.com, 9098, IAM auth (별도 JAAS 설정 필요)
Azure Event Hubs (Kafka API)host = <ns>.servicebus.windows.net, 9093 + SASL_SSL + PLAIN + Connection String JAAS
신규 consumer group + earliest기존 데이터 모두 받기 위해 AUTO OFFSET=earliest 선택
group offset 관리 분리다른 GROUP ID 로 등록 — 같은 topic 도 group 별로 별도 offset 관리됨

write (producer.send)

태그 페이지 또는 REST API 로 값을 쓰면 producer 가 해당 topic 에 record 로 전송됩니다 (async, fire-and-forget). value 의 String 이 record value 로 전송됩니다 — JSON / 평문 모두 가능.

plc_address = factory/line1/cmd
value = ON
→ Kafka producer.send: topic="factory/line1/cmd" value="ON"
plc_address = factory/line1/cmd:temperature
value = 25.3
→ Kafka producer.send: topic="factory/line1/cmd" value="25.3" (콜론 뒤는 read decoder 정보, write 는 topic 만 사용)

흔한 문제와 해결

증상원인해결
값이 안 들어옴broker 가 record 안 보냄kafka-console-consumer.sh --bootstrap-server <host>:9092 --topic <t> --from-beginning 으로 직접 받아보고 확인
[KAFKA] connect 실패: TimeoutExceptionbootstrap 서버 도달 불가telnet <host> 9092 검증. advertised.listeners 가 외부 IP 인지 확인 (broker 측)
[KAFKA] poll error: Authentication failedSASL 인증 실패sasl-mechanism / sasl-jaas-config 재확인. password 특수문자 escape
[KAFKA] poll error: SslHandshakeTLS truststore 누락 / 만료broker 인증서 신뢰 가능한 CA 발급 또는 시스템 truststore 추가
신규 group 가 데이터 못 받음auto-offset=latest 인데 신규 record 가 없음일회성으로 earliest 로 변경 또는 publish 가 발생할 때까지 대기

더 자세한 기술 문서