You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

咨询Confluent Schema Registry本地缓存的查询及管理方法

Inspecting Confluent Avro Serializer/Deserializer Local Cache

Great question—when pushing millions of Kafka messages with Confluent Schema Registry (SR), ensuring the local cache in AvroSerializer/AvroDeserializer is working as expected is critical to avoid unnecessary SR calls. Let me break down how you can check its contents, size, and overall behavior:

1. Understand the Cache Implementation

First, note that the cache is managed by the CachedSchemaRegistryClient (the default client used by Avro serializers/deserializers). It uses two core caches under the hood:

  • A cache mapping schema IDs to Schema objects (for deserialization, when fetching schemas by ID)
  • A cache mapping subjects to schema ID/Schema pairs (for serialization, when registering or looking up schemas by subject)

By default, this cache has a maximum size of 1000 entries and no expiration time.

2. Inspect Cache Contents/Size in Code (Debugging)

If you’re in a development environment and want to directly view the cache, you can use reflection to access the private cache fields in CachedSchemaRegistryClient (there’s no public API for this out of the box):

Example Code (Java)

import org.apache.avro.Schema;
import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient;
import com.google.common.cache.LoadingCache;
import java.lang.reflect.Field;
import java.util.Map;
import java.util.Collections;

public class CacheInspector {
    public static void main(String[] args) throws Exception {
        // Initialize your Schema Registry client
        CachedSchemaRegistryClient client = new CachedSchemaRegistryClient(
                "http://your-sr-url:8081",
                1000 // Default max cache size
        );

        // Access the schema ID -> Schema cache
        Field schemaCacheField = CachedSchemaRegistryClient.class.getDeclaredField("schemaCache");
        schemaCacheField.setAccessible(true);
        LoadingCache<Integer, Schema> schemaCache = (LoadingCache<Integer, Schema>) schemaCacheField.get(client);

        // Print cache size and entries
        System.out.println("Schema ID Cache Size: " + schemaCache.size());
        for (Map.Entry<Integer, Schema> entry : schemaCache.asMap().entrySet()) {
            System.out.printf("Schema ID: %d | Schema Name: %s%n", entry.getKey(), entry.getValue().getName());
        }

        // Access the subject -> schema ID cache (for serialization lookups)
        Field subjectCacheField = CachedSchemaRegistryClient.class.getDeclaredField("subjectCache");
        subjectCacheField.setAccessible(true);
        LoadingCache<String, Integer> subjectCache = (LoadingCache<String, Integer>) subjectCacheField.get(client);

        System.out.println("\nSubject Cache Size: " + subjectCache.size());
        for (Map.Entry<String, Integer> entry : subjectCache.asMap().entrySet()) {
            System.out.printf("Subject: %s | Schema ID: %d%n", entry.getKey(), entry.getValue());
        }
    }
}

Note: Reflection works great for debugging, but avoid using it in production code—it can break if Confluent changes the internal class structure.

3. Monitor Cache Performance in Production (JMX)

For production environments, use JMX metrics to track cache behavior without modifying code. Confluent’s Schema Registry client exposes these key metrics:

  • kafka.schema.registry.client.cache.hit.ratio: The percentage of SR requests served from the cache (aim for >95% for optimal performance)
  • kafka.schema.registry.client.cache.size: Current number of entries in the cache
  • kafka.schema.registry.client.cache.miss.count: Total number of cache misses (indicates how often the client had to call the SR)

You can view these metrics using tools like:

  • JConsole or VisualVM (connect directly to your Kafka producer/consumer JVM)
  • Prometheus + Grafana (if you’ve set up JMX Exporter to scrape metrics)

4. Adjust Cache Configuration

If you find the cache is too small (leading to frequent misses), you can increase its maximum size when initializing the CachedSchemaRegistryClient:

Example: Custom Cache Size

int maxCacheSize = 5000; // Increase to handle more schema entries
CachedSchemaRegistryClient client = new CachedSchemaRegistryClient(
        "http://your-sr-url:8081",
        maxCacheSize,
        Collections.emptyMap(),
        1000 // Initial cache capacity
);

If you’re using Spring Kafka, set this via properties:

spring.kafka.producer.properties.schema.registry.client.cache.capacity=5000
spring.kafka.consumer.properties.schema.registry.client.cache.capacity=5000

内容的提问来源于stack exchange,提问作者user954156

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:06:43