diff --git a/test/kafka/kafka-client-loadtest/docker-compose.yml b/test/kafka/kafka-client-loadtest/docker-compose.yml index 00b3264eb..fa88117e9 100644 --- a/test/kafka/kafka-client-loadtest/docker-compose.yml +++ b/test/kafka/kafka-client-loadtest/docker-compose.yml @@ -62,6 +62,8 @@ services: SCHEMA_REGISTRY_KAFKASTORE_WRITE_TIMEOUT_MS: "60000" SCHEMA_REGISTRY_KAFKASTORE_INIT_RETRY_BACKOFF_MS: "5000" SCHEMA_REGISTRY_KAFKASTORE_CONSUMER_AUTO_OFFSET_RESET: "earliest" + # Enable debug logging for Kafka clients to trace consumer seek/fetch operations + SCHEMA_REGISTRY_LOG4J_LOGGERS: "org.apache.kafka.clients=DEBUG,org.apache.kafka.common=DEBUG,io.confluent.kafka.schemaregistry.storage=DEBUG" healthcheck: test: ["CMD", "curl", "-f", "http://localhost:8081/subjects"] interval: 15s diff --git a/weed/mq/kafka/protocol/handler.go b/weed/mq/kafka/protocol/handler.go index 509ed1330..2f967c272 100644 --- a/weed/mq/kafka/protocol/handler.go +++ b/weed/mq/kafka/protocol/handler.go @@ -1716,6 +1716,16 @@ func (h *Handler) HandleMetadataV3V4(correlationID uint32, requestBody []byte) ( response := buf.Bytes() + // Detailed logging for Metadata response + maxDisplay := len(response) + if maxDisplay > 50 { + maxDisplay = 50 + } + glog.Warningf("🟡 Metadata v3/v4 FINAL RESPONSE: size=%d bytes, first50bytes=%v", len(response), response[:maxDisplay]) + if len(response) > 100 { + glog.Warningf("🟡 Metadata v3/v4 RESPONSE HEX (first 100 bytes): %x", response[:100]) + } + return response, nil } @@ -1882,6 +1892,16 @@ func (h *Handler) handleMetadataV5ToV8(correlationID uint32, requestBody []byte, response := buf.Bytes() + // Detailed logging for Metadata response + maxDisplay := len(response) + if maxDisplay > 50 { + maxDisplay = 50 + } + glog.Warningf("🟡 Metadata v%d FINAL RESPONSE: size=%d bytes, first50bytes=%v", apiVersion, len(response), response[:maxDisplay]) + if len(response) > 100 { + glog.Warningf("🟡 Metadata v%d RESPONSE HEX (first 100 bytes): %x", apiVersion, response[:100]) + } + return response, nil }