클러스터 생성
kind: Cluster
apiVersion: kind.x-k8s.io/v1alpha4
nodes:
- role: control-plane
image: kindest/node:v1.33.1
extraPortMappings: // 외부에서 접근 가능하도록 30080 을 열어줘야 함
- containerPort: 30080
hostPort: 30080
protocol: TCP
- containerPort: 30081
hostPort: 30081
protocol: TCP
helm repo 추가 & repo update
helm repo **add** strimzi https://strimzi.io/charts/
helm repo update
Strimzi Cluster Operator 설치
helm **install strimzi-kafka-operator** strimzi/strimzi-kafka-**operator** --namespace strimzi-kafka \ # 네임스페이스 지정
# --version 0.8.2 # 특정 버전 지정
kafka 생성
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaNodePool
metadata:
name: dual-role
labels:
strimzi.io/cluster: my-cluster
spec:
replicas: 1
roles:
- controller
- broker
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 100Gi
deleteClaim: false
kraftMetadata: shared
---
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
name: my-cluster
annotations:
strimzi.io/node-pools: enabled
strimzi.io/kraft: enabled
spec:
kafka:
version: 4.0.0
metadataVersion: 4.0-IV3
listeners:
- name: plain
port: 9092
type: internal
tls: false
- name: tls
port: 9093
type: internal
tls: true
- name: external
port: 9094
type: nodeport # 또는 loadbalancer
tls: false
configuration:
brokers:
- broker: 0
advertisedHost: 172.30.1.71
advertisedPort: 30092
config:
offsets.topic.replication.factor: 1
transaction.state.log.replication.factor: 1
transaction.state.log.min.isr: 1
default.replication.factor: 1
min.insync.replicas: 1
entityOperator:
topicOperator: {}
userOperator: {}
kafka-topic 생성
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
name: my-topic
labels:
strimzi.io/cluster: my-cluster
spec:
partitions: 1
replicas: 1
config:
retention.ms: 7200000
segment.bytes: 1073741824
kafka-bridge 생성
# 외부 통신이 가능하도록 변경한 후
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaBridge
metadata:
name: my-bridge
spec:
replicas: 1
bootstrapServers: my-cluster-kafka-bootstrap:9092
http:
port: 8080
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaBridge
metadata:
name: my-bridge
spec:
replicas: 1
# boostrapServers 값이 중요, 접근이 가능한 bootstrapServers 를 지정해야 함
# 서비스 이름을 지정해도 되는데, 구성할 때 접근이 불가능했음
bootstrapServers: my-cluster-kafka-bootstrap.strimzi-operator.svc:9092
http:
port: 8080
명령 - kubectl get kafkabridge my-bridge -n strimzi-operator -o yaml | grep bootstrapServers
결과 - {"apiVersion":"kafka.strimzi.io/v1beta2","kind":"KafkaBridge","metadata":{"annotations":{},"name":"my-bridge","namespace":"strimzi-operator"},"spec":{"bootstrapServers":"172.30.1.71:30092","http":{"port":8080},"replicas":1}}
bootstrapServers: 172.30.1.71:30092
master@master-T8-Plus:~/kafka$
apiVersion: v1
kind: Service
metadata:
name: my-bridge-nodeport
namespace: strimzi-operator
spec:
type: NodePort
selector:
strimzi.io/kind: KafkaBridge
ports:
- port: 8080
targetPort: 8080
nodePort: 30080 # 외부에서 30080 포트로 접근하면 KafkaBridge - 8080 포트로 전달
protocol: TCP
연동 테스트
내부 Pod 접근해서 테스트
# Producer 실행해서 메시지 전송
kubectl exec -it my-cluster-dual-role-0 -n strimzi-operator -- \
bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic my-topic
# topic (my-topic) 에서 메시지 확인
kubectl exec -it my-cluster-dual-role-0 -n strimzi-kafka -- \
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic my-topic --from-beginning
kafka bridge 로 테스트 (클러스터를 외부 설정을 안해놔서 포트포워딩으로 진행함)
# localhost:30080 접근 시 8080 으로 포워딩
# 외부 접근이 가능하면 포트포워딩 필요 없음
kubectl port-forward svc/my-bridge-nodeport -n strimzi-operator 30080:8080
# 메시지 전송
curl -X POST "http://localhost:30080/topics/my-topic" \
-H "Content-Type: application/vnd.kafka.json.v2+json" \
-d '{"records":[{"value":{"a":"a"}}]}'
# Consumer 인스턴스 생성
curl -X POST http://localhost:30080/consumers/my-consumer \
-H "Content-Type: application/vnd.kafka.v2+json" \
-d '{
"name": "my-consumer-instance",
"format": "json",
"auto.offset.reset": "earliest"
}'
# Topic 구독 -> my-topic 을 구동하도록 함
curl -X POST http://localhost:30080/consumers/my-consumer/instances/my-consumer-instance/subscription \
-H "Content-Type: application/vnd.kafka.v2+json" \
-d '{
"topics": ["my-topic"]
}'
# 메시지 소비
curl -X GET http://localhost:30080/consumers/my-consumer/instances/my-consumer-instance/records \
-H "Accept: application/vnd.kafka.json.v2+json"