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 ID | consumer group id (옵션) | plant-floor-A, edge-line1 |
| SECURITY | 보안 프로토콜 | PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL |
| SASL MECH | SASL 메커니즘 (옵션) | PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 |
| SASL JAAS | SASL 인증 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 입니다.
| 모드 | 형식 | 동작 |
|---|---|---|
| SCALAR | sensors/temp 또는 sensors/temp.value | 메시지 전체를 String 으로. 메시지가 JSON object 면 raw 폴백. |
| KEY | sensors/temp:temperature | top-level JSON key 의 값 (예: {"temperature":25.3,"humidity":60} → 25.3) |
| PATH | sensors/temp:$.data.tags.T1 | JSON Pointer 동적 evaluate (중첩 key 지원) |
| RAW | sensors/temp:_raw_ | 마지막 record value 전체 (디버깅) |
read 의 첫 호출은 lazy subscribe (빈 문자열 반환). 다음 polling cycle 부터 캐시값이 들어옵니다.
자주 쓰는 사례
| 사례 | 어떻게 |
|---|---|
| 사내 Kafka cluster | host = bootstrap broker IP, 9092 (평문) |
| Confluent Cloud | host = 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 실패: TimeoutException | bootstrap 서버 도달 불가 | telnet <host> 9092 검증. advertised.listeners 가 외부 IP 인지 확인 (broker 측) |
[KAFKA] poll error: Authentication failed | SASL 인증 실패 | sasl-mechanism / sasl-jaas-config 재확인. password 특수문자 escape |
[KAFKA] poll error: SslHandshake | TLS truststore 누락 / 만료 | broker 인증서 신뢰 가능한 CA 발급 또는 시스템 truststore 추가 |
| 신규 group 가 데이터 못 받음 | auto-offset=latest 인데 신규 record 가 없음 | 일회성으로 earliest 로 변경 또는 publish 가 발생할 때까지 대기 |