레이블이 kafka인 게시물을 표시합니다. 모든 게시물 표시
레이블이 kafka인 게시물을 표시합니다. 모든 게시물 표시

kafka invalid timestamp error by pidstat

k8s 환경에서 프로세스등의 로그를 수집(collector pod) > json stdout(/var/log/pods/.../0.log) > filebeat > kafka cluster 로 흐름이 구축되어 있다.
그런데 자정(0시)에 가끔 collector CPU 90% 이 넘어가고 해소가 되지 않는다.

kafka(4.0.2) server.log 를 보면 다음과 같이 invalid timestamp 로 거부가 됐다.
kafka 디폴트로 메시지 1시간까지 허용하고 넘어가면 거부된다.
org.apache.kafka.common.errors.InvalidTimestampException: One or more records have been rejected due to invalid timestamp

collector 는 pidstat 커맨드로 수집을 하는데 pidstat 는 Time=HH:MM:SS 으로 날짜가 출력되지 않는다.
# -h : 한 줄 가로 출력, 평균 행 없음 (파싱용)
# -u : CPU
# -r : 메모리
# -d : 디스크 I/O
# -l : Command 에 전체 커맨드라인
# -I : CPU 사용률을 전체 코어 수로 나눔 (SMP 정규화)
# -T : TASK + CHILD 통계 모두
# -p : 대상 PID (discover 가 찾은 프로세스들)
pidstat -hurdl -I -T ALL -p <PIDs> <interval> <count>

# Time        UID       PID    %usr %system  %guest   %wait    %CPU   CPU  minflt/s  majflt/s     VSZ     RSS   %MEM   kB_rd/s   kB_wr/s kB_ccwr/s iodelay  Command
05:04:08        0       630    0.00    0.00    0.00    0.00    0.00     9      0.00      0.00    2228    1024   0.01      0.00      0.00      0.00       0  sleep 300

이 Time을 collector 에서 파싱하고 now() 의 날짜부분과 합쳐서 @timestamp 필드값을 설정한다.

# 파싱 과정
# local_to_utc() python 로직
dt = datetime.strptime(local_timestamp, "%H:%M:%S")
return datetime.now().replace(hour=dt.hour, minute=dt.minute, second=dt.second).astimezone(tz=timezone.utc)

# timestamp 필드 설정
obj["@timestamp"] = local_to_utc(obj.pop("Time")).strftime("%Y-%m-%dT%H:%M:%SZ")

# 문제가 되는 상황 예시
pidstat 출력 23:59:59.6 (8/15) -> Time="23:59:59"
파서 처리 시점 00:00:00.2 (8/16) -> now().date = 8/16 날짜 붙여서 -> timestamp = 2026-0816T23:59:59Z
실제 측정 시각 8/15 23:59:59 대비 +24h이 지난 상태의 timestamp 가 된다.
+24h 메시지는 kafka 에서 거부하는데, collector 는 계속 재시도를 해 cpu 리소스가 높게 유지 된 것이다.

# 해결 시도
# pidstat 와 1초 정도밖에 차이 안나서 이렇게 파싱시 now 를 timestamp 로 사용할 수 있겠지만
# pidstat 이 한 주기에 프로세스 n개를 측정하고 있어 원래는 n개의 라인이 전부 Time 이 같아야 한다. 매번 now 로 조금씩 달라지면 배치 동일성이 깨진다.
obj["@timestamp"] = datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")

# 해결 방법1
# +24(12시간 차이 정도로 하자)시간이 지난 경우 -1day 한다.
now = datetime.now()
ts = now.replace(hour=dt.hour, minute=dt.minute, second=dt.second, microsecond=0)
if ts - now > timedelta(hours=12):
    ts -= timedelta(days=1)
    logger.info(f"midnight rollover: Time={local_timestamp} -> {ts}")
return ts.astimezone(tz=timezone.utc)

# 해결 방법2
# pidstat 사용 안 하고 k8s 환경이면 이미 cAdvisor(ContainerAdvisor, kubelet 내장, 실행 중인 컨테이너들의 리소스 사용량과 성능 메트릭을 수집,처리,노출하는 도구)로 prometheus 에 메트릭이 있으니 이 값을 조회해서 사용한다.
# 메트릭들은 이미 같은 timestamp 로 저장되어 있다.
# 예시
%CPU
sum by (namespace, pod, container) (
 rate(container_cpu_usage_seconds_total{container!="",container!="POD",namespace=~"ysoftman.*"}[5m])
) * 100

%usr
sum by (namespace, pod, container) (
 rate(container_cpu_user_seconds_total{container!="",container!="POD",namespace=~"ysoftman.*"}[5m])
) * 100

%system
sum by (namespace, pod, container) (
 rate(container_cpu_system_seconds_total{container!="",container!="POD",namespace=~"ysoftman.*"}[5m])
) * 100

RSS - pidstat RSS 단위는 KB, cAdvisor는 byte
sum by (namespace, pod, container) (
 container_memory_rss{container!="",container!="POD",namespace=~"ysoftman.*"}
) / 1024

# 현재시각을 고정하고 이 값으로 prometheus 메트릭을 조회 하고 그대로 @timestamp 로 설정한다.
at = datetime.now(timezone.utc).timestamp()
ts = datetime.fromtimestamp(at, timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")

# 다음과 같이 만들수 있다.
{"@timestamp": "2026-08-16T05:23:44Z", "%CPU": "0.88", "%usr": "35.59", "%system": "6.81", "RSS": "571528", "container": "ysoftman1-container", "pod": "ysoftman1-aaa"}

# 어쩌다 발생해서 신경쓰였는데 원인 파악 돼서 개비스콘 짤로 남긴다.ㅎ

filebeat kafka out of acceptable range

log 수집용 filebeat pod 하나가 계속 높은 CPU 사용량을 보인다.
로그를 보면
{"log.level":"error","@timestamp":"2026-04-07T08:12:32.206+0900","log.logger":"kafka","log.origin":{"function":"github.com/elastic/beats/v7/libbeat/outputs/kafka.(*msgRef).dec","file.name":"kafka/client.go","file.line":446},"message":"Kafka publish failed with: kafka server: The timestamp of the message is out of acceptable range","service.name":"filebeat","ecs.version":"1.6.0"}

pod 를 재시작해도 시간이 좀 지나면 다시 CPU 사용량이 90%를 넘어간다.

kafka publish 실패 -> 재시도 무한 반복 -> CPU 급증

broker 레벨 timestamp 설정 확인
kafka-configs --bootstrap-server 카프카호스트:9092 --entity-type brokers --describe --all | rg -i timestamp

log.message.timestamp.type=CreateTime (체크 대상 타입)
log.message.timestamp.after.max.ms=3600000 (1시간, 메시지 시간이 이 시간보다 1시간 초과시 거부)
log.message.timestamp.before.max.ms=9223372036854775807 (사실 무제한, 과거 메시지 모두 허용)

메시지 timestamp 가 before / after 를 벗어난 경우 kafka broker가 timestamp 범위 초과로 메시지를 거부한다.
kafka timestamp 설정은 이상이 없어 보인다.

filebeat.yml 설정에 max_retries가 명시되어 있지 않아 디폴트 3으로 사용된다.
output.kafka:
    enabled: true
    hosts: "${kafka_hosts}"
    topic: "%{[@metadata][topic]}"
    partition.round_robin:
      reachable_only: false
      group_events: 1

max_retries 회수가 제한이 있는데도 CPU가 안 떨어진다는 건 처리해야 할 오래된 로그 라인 자체가 엄청나게 많다는뜻일 수 있다.
registry 초기화 후 모든 로그 파일을 처음부터 다시 읽으면서
이벤트마다 전송 시도 -> 실패 -> retry 3번 (backoff 포함) -> drop -> 다음 이벤트 -> 또 실패
이 사이클이 수십만 건 이상 반복되고 있을 수 있습니다. 시간이 충분히 지나면 결국 CPU는 내려가겠지만, 로그 양이 많으면 오래 걸릴 수 있습니다.

당장 해결을 위해선 pod 재시작 한다.
pod 재시작 후 에도 같은 현상이면 container 상에서 registry(active.dat, log.json, 12345.json)을 지우고 재시작하자. 이때 meta.json 을 지우면 pod 재시작시 crash 가 발생하니 유지해야 된다.
실수로 삭제했다면 meta.json 는 이런 내용으로 다시 생성하면 된다.
{"version":"1"}

filebeat 수집 대상이 *.log* 로 되어 있어 로그 로테이션으로 log 파일들도 매칭이 되고 있는게 근본적인 원인으로 보여 *.log* -> *.log 패턴으로 수정했다.

argo-cd cilium identity

k8s: v1.26.4
argo-cd: v2.12.3+6b9cd82
에서 strimzi kafka operator 로 kafka cluster 를 설치후
argocd 에서 보면 다음과 같이 cilium identity 에 out of sync 로 보인다.

실제 CiliumIdentity 를 보면 정상적으로 리소스가 생성돼 있다.
관련 이슈가 cilium troubleshooting 이 있다.
다음과 같이 argocd 에서 CiliumIdentity 리소스를 모니터링 대상에서 제외시킨다.
kubectl edit configmap argocd-cm -n argocd
resource.exclusions: |
 - apiGroups:
     - cilium.io
   kinds:
     - CiliumIdentity
   clusters:
     - "*"

# helm 으로 반영시
# values 는 다음과 같이 수정하자.
configs:
  cm:
    resource.exclusions: |
     - apiGroups:
         - cilium.io
       kinds:
         - CiliumIdentity
       clusters:
         - "*"

# 반영
helm upgrade argocd argo/argo-cd \
--namespace argocd \
--values values-dev.yaml

# argocd 3 부터 CiliumIdentity 가 기본 제외되어 별도 설정을 하지 않아도 된다.

#####

strimzi-ysoftman-kafka-cluster-ysoftman-kafka-cluster-kafka-init
라는 이름의 ClusterRoleBinding(crb) 도 실제 리소스가 있지만 argocd 에서는 out of sync 로 표시된다.
관련 논의가 있지만 아직 해결방법은 없는것 같다.

위 argocd cm > resource.exclusions 적용해도 out of sync 로 나온다.
차선책으로 ...kafka-init crb 에 다음을 추가하면 kafka 전체 sync 에서는 제외시켜 kafka-cluster 앱으로는 synced 로 만들 수는 있다.
metadata:
  annotations:
    argocd.argoproj.io/compare-options: IgnoreExtraneous

# argocd 3.x 로 버전업하니 확인결과 해결되었다.
# strimzi-xxx-kafka-cluster-xxx-init (crb)는 존재하지만 argocd 싱크대상에서 제외된다.

strimzi kafka nodeport ingress

# strimzi operator 로 k8s 에 kafka cluster 를 구성한 경우
# 클러스터들이 svc 로컬 호스트 사용으로 k8s 클러스터 외부에서 kafka 9092포트(bootstrap/broker)로 접속이 안된다.

# 우선 kafka 설치가 되어 있어야 테스트할 수 있다.
# /opt/homebrew/opt/kafka/bin 사용할 수 있는 커맨드 스크립트들이 생성된다.
brew install kafka kcat

# nodeport 생성하기
# kafka 리소스 > spec > kafka > listeners 에 다음과 설정을 추가하면 nodeport 가 생성된다.
# service/pod 에 9094 nodeport 설정이 추가된다.
- name: external1 # ^[a-z0-9]{1,11}$' 이고 유니크해야 한다.
  port: 9094
  type: nodeport
  tls: false
  configuration:
    bootstrap:
      nodePort: 32100
    brokers:
    - broker: 0
      nodePort: 32000
    - broker: 1
      nodePort: 32001
    - broker: 2
      nodePort: 32002

# nodeport 접속 확인
# 토픽으로 메시지 생성
/opt/homebrew/opt/kafka/bin/kafka-console-producer \
--broker-list ysoftman-node1:32100 \
--topic test

# 토픽으로 들오는 메시지 확인
/opt/homebrew/opt/kafka/bin/kafka-console-consumer \
--bootstrap-server ysoftman-node1:32100 \
--topic test \
--from-beginning
# 또는
kcat -b ysoftman-node1:32100 -t test

#####

# ingress 생성하기
# ingress 는 http 프로토콜을 사용하지만 kafka 는 tcp 프로토콜을 사용한다.
# 따라서 nginx ingress > ssl-passthrough 기능을 사용해 서비스 tcp 로 바로 연결되는 방식을 사용해야 한다.
# kafka 리소스 > spec > kafka > listeners 에 다음과 설정을 추가하면 ingress 가 생성된다.
# service/pod 에 9096 포트 설정이 추가된다.
- name: external1 # ^[a-z0-9]{1,11}$' 이고 유니크해야 한다.
  port: 9096
  tls: true # Ingress type listener and requires enabled TLS encryption
  type: ingress
  configuration:
    bootstrap:
      host: ysoftman-bootstrap.ysoftman.abc
    brokers:
    - broker: 0
      host: ysoftman-0.ysoftman.abc
    - broker: 1
      host: ysoftman-1.ysoftman.abc
    - broker: 2
      host: ysoftman-2.ysoftman.abc
    class: nginx # kubectl get ingressclass 로 사용 가능한 클래스 이름 파악

# 잠시 후 생성된 인그레스 중 하나를 보면 다음과 같다.
# tls 에 별도의 secretName 이 없다.
# 대신 ssl-passthrough 활성화한다.
# nginx-ingress-controller daemonset(또는 deployment) 에 --enable-ssl-passthrough 설정을 적용해야 ingress ssl-passthrough 이 동작한다.
spec:
  template:
    spec:
      containers:
      - args:
        - /nginx-ingress-controller
        - --enable-ssl-passthrough=true

# kafka 서버(broker)에서 https 를 받고 tls 인증을 처리하게 된다.
# 참고
metadata:
  annotations:
    ingress.kubernetes.io/ssl-passthrough: "true"
    nginx.ingress.kubernetes.io/backend-protocol: HTTPS
    nginx.ingress.kubernetes.io/ssl-passthrough: "true"
... 생략 ...
spec:
  tls:
  - hosts:
    - ysoftman-bootstrap.ysoftman.abc

# 그리고 client, cluster 등의 이름으로 secret 도 생성이 된다.
# 이중 client secret 를 .crt 파일로 다음과 같이 저장한다.
kubectl get secret ysoftman-kafka-cluster-cluster-ca-cert -o jsonpath='{.data.ca\.crt}' | base64 -d > ca.crt

# kafka 커맨드에서 사용할 truststore.jks 파일 생성
keytool -import -trustcacerts -alias root -file ca.crt -keystore truststore.jks -storepass password -noprompt

# kafka 클러스터에 접속해서 producing 해보기
/opt/homebrew/opt/kafka/bin/kafka-console-producer --broker-list ysoftman-bootstrap.ysoftman.abc:443 --producer-property security.protocol=SSL --producer-property ssl.truststore.password=password --producer-property ssl.truststore.location=./truststore.jks --topic test

# kafka 클러스터에 접속해서 consume 해보기
/opt/homebrew/opt/kafka/bin/kafka-console-consumer --bootstrap-server ysoftman-bootstrap.ysoftman.abc:443 --topic test

# 인증서 확인
openssl s_client -connect ysoftman-bootstrap.ysoftman.abc:443 \
-servername ysoftman-bootstrap.ysoftman.abc \
-showcerts

# 만약 다음과 같은 ssl 실패 에러가 발생한다면
failed authentication due to: SSL handshake failed

# ssl 디버깅 정보를 보자
export KAFKA_OPTS="-Djavax.net.debug=ssl"

#####

# strimzi operator 로 kafka 를 설치한 경우 broker, controller pod 들은
# strimzipodset(statefulset 과 비슷) 이라는 커스텀 리소스로 관리된다.
# broker pod 1개를 수동 삭제했는데 pod 가 새로 올라 올때 다른 pod 들과 연결 에러가 발생한다.
# 테스트해본 결과 strimzipodset broker, controller 모두 삭제해서 재시작하도록 하면 된다.
managed-kafka-cluster-broker
managed-kafka-cluster-controller

redis, kafka 유입 확인

# redis 유입 확인
# 참고 https://redis.io/commands/monitor
# redis 서버라면
redis-clit monitor

# 아니라면 텔넷으로 확인
telnet ysoftman-redis 6379

# 접속 후 monitor 커맨드로 유입 확인
monitor

# 종료
ctrl + ] 후 quit

# 또는 redis-cli 명령 사용
redis-cli -h ysoftman-redis -p 6379 monitor

#####

# kafkacat(kcat) 으로 유입 확인
# https://github.com/edenhill/kafkacat
# brew install kafkacat

# 1.7.0 부터 kcat 으로 이름이 변경됨.
# 맥에서 설치
brew install kcat

# -b brocker(kafaka 서버), -t topic(쿼리할 토픽)
kcat -b ysoftman_host:9092 -t my_topic

# -L 메타데이터(브로커,토픽) 리스트 보기
kcat -L -b ysoftman_host:9092

Prev