Windows Server环境下Kafka Broker容器报Invalid Argument错误排查
问题
我通过Docker Compose部署了Kafka容器集群,具体配置如下:
kafka: image: "confluentinc/cp-kafka" container_name: kafka restart: always depends_on: - zookeeper ports: - '9092:9092' environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 100 CONFLUENT_METRICS_ENABLE: 'false' networks: - ddn deploy: resources: limits: memory: 2G init-kafka: image: confluentinc/cp-kafka depends_on: - kafka entrypoint: [ '/bin/sh', '-c' ] command: | " # blocks until kafka is reachable kafka-topics --bootstrap-server kafka:29092 --list echo -e 'Creating kafka topics' kafka-topics --bootstrap-server kafka:29092 --create --if-not-exists --topic EVENT_TYPE1_NOTIFY --replication-factor 1 --partitions 1 echo -e 'Successfully created the following topics:' kafka-topics --bootstrap-server kafka:29092 --list " networks: - ddn kafka-rest: container_name: kafka-rest image: confluentinc/cp-kafka-rest hostname: kafka-rest restart: always depends_on: - init-kafka ports: - "8082:8082" environment: KAFKA_REST_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_REST_HOST_NAME: kafka-rest KAFKA_REST_LISTENERS: http://kafka-rest:8082 KAFKA_REST_BOOTSTRAP_SERVERS: kafka:29092 healthcheck: test: ["CMD", "curl", "-f", "http://kafka-rest:8082/topics/EVENT_TYPE1_NOTIFY "] interval: 5s timeout: 1s retries: 12 networks: - ddn networks: ddn:
但在基于Windows Server的集成测试代理中运行时,出现如下错误:
[2023-01-03 13:11:26,453] ERROR [KafkaApi-1] Error when handling request: clientId=1, correlationId=5, api=LEADER_AND_ISR, version=5, body=LeaderAndIsrRequestData(controllerId=1, controllerEpoch=1, brokerEpoch=25, type=0, ungroupedPartitionStates=[], topicStates=[LeaderAndIsrTopicState(topicName='EVENT_TYPE1_NOTIFY', topicId=MU_GDxwKTkywLg9wPNCrsg, partitionStates=[LeaderAndIsrPartitionState(topicName='EVENT_TYPE1_NOTIFY', partitionIndex=0, controllerEpoch=1, leader=1, leaderEpoch=0, isr=[1], zkVersion=0, replicas=[1], addingReplicas=[], removingReplicas=[], isNew=true)])], liveLeaders=[LeaderAndIsrLiveLeader(brokerId=1, hostName='kafka', port=9092)]) (kafka.server.RequestHandlerHelper) java.io.IOException: Invalid argument at java.base/sun.nio.ch.FileChannelImpl.map0(Native Method) at java.base/sun.nio.ch.FileChannelImpl.map(Unknown Source) at kafka.log.AbstractIndex.<init>(AbstractIndex.scala:124) at kafka.log.OffsetIndex.<init>(OffsetIndex.scala:54) at kafka.log.LazyIndex$.$anonfun$forOffset$1(LazyIndex.scala:106) at kafka.log.LazyIndex.$anonfun$get$1(LazyIndex.scala:63) at kafka.log.LazyIndex.get(LazyIndex.scala:60) at kafka.log.LogSegment.offsetIndex(LogSegment.scala:64) at kafka.log.LogSegment.readNextOffset(LogSegment.scala:456) at kafka.log.Log.$anonfun$recoverLog$6(Log.scala:944) at scala.runtime.java8.JFunction0$mcJ$sp.apply(JFunction0$mcJ$sp.scala:17) at scala.Option.getOrElse(Option.scala:201) at kafka.log.Log.recoverLog(Log.scala:944) at kafka.log.Log.$anonfun$loadSegments$3(Log.scala:824) at scala.runtime.java8.JFunction0$mcJ$sp.apply(JFunction0$mcJ$sp.scala:17) at kafka.log.Log.retryOnOffsetOverflow(Log.scala:2494) at kafka.log.Log.loadSegments(Log.scala:824) at kafka.log.Log.<init>(Log.scala:328) at kafka.log.Log$.apply(Log.scala:2630) at kafka.log.LogManager.$anonfun$getOrCreateLog$1(LogManager.scala:830) at scala.Option.getOrElse(Option.scala:201) at kafka.log.LogManager.getOrCreateLog(LogManager.scala:783) at kafka.cluster.Partition.createLog(Partition.scala:344) at kafka.cluster.Partition.createLogIfNotExists(Partition.scala:324) at kafka.cluster.Partition.$anonfun$makeLeader$1(Partition.scala:563) at kafka.cluster.Partition.makeLeader(Partition.scala:547) at kafka.server.ReplicaManager.$anonfun$makeLeaders$5(ReplicaManager.scala:1568) at kafka.utils.Implicits$MapExtensionMethods$.$anonfun$forKeyValue$1(Implicits.scala:62) at scala.collection.mutable.HashMap$Node.foreachEntry(HashMap.scala:633) at scala.collection.mutable.HashMap.foreachEntry(HashMap.scala:499) at kafka.server.ReplicaManager.makeLeaders(ReplicaManager.scala:1566) at kafka.server.ReplicaManager.becomeLeaderOrFollower(ReplicaManager.scala:1411) at kafka.server.KafkaApis.handleLeaderAndIsrRequest(KafkaApis.scala:258) at kafka.server.KafkaApis.handle(KafkaApis.scala:171) at kafka.server.KafkaRequestHandler.run(KafkaRequestHandler.scala:74) at java.base/java.lang.Thread.run(Unknown Source)
请问是否存在配置缺失?
解决方案
这个错误并非配置缺失,而是Windows Server的文件系统与Kafka默认的索引文件内存映射机制不兼容导致的。以下是具体修复步骤:
- 添加Kafka索引类型配置
在kafka服务的environment中加入如下配置,强制Kafka使用文件模式而非内存映射模式处理索引文件:
KAFKA_LOG_INDEX_TYPE: "file"
- 挂载本地卷存储Kafka数据
Windows环境下,容器匿名卷的文件系统可能存在兼容性问题,需为Kafka添加本地卷挂载,将日志数据持久化到Windows主机目录:
在kafka服务中新增volumes节点:
volumes: - ./kafka-data:/var/lib/kafka/data
提前在主机创建kafka-data目录,并确保Docker拥有该目录的读写权限。
- 修正init-kafka容器的命令格式
原command中的嵌套引号会导致shell解析异常,修改为:
command: | kafka-topics --bootstrap-server kafka:29092 --list echo -e 'Creating kafka topics' kafka-topics --bootstrap-server kafka:29092 --create --if-not-exists --topic EVENT_TYPE1_NOTIFY --replication-factor 1 --partitions 1 echo -e 'Successfully created the following topics:' kafka-topics --bootstrap-server kafka:29092 --list
去掉外层双引号,避免命令执行失败。
- 验证网络与端口权限
确保Windows Server防火墙允许9092、8082端口的入站/出站流量,同时Docker网络ddn能正常连通所有服务容器。
内容的提问来源于stack exchange,提问作者Or S
相关产品推荐
相关产品推荐

