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

Java读写MSK主题:键值为AVRO并使用双Glue Schema Registry方法

解决Glue Schema Registry中键值分离注册表/Schema名称的配置问题

核心解决方案

Glue Schema Registry(GSR)的Kafka序列化/反序列化器支持通过**key.和value.前缀区分键与值的配置参数**——虽然AWSSchemaRegistryConstants只定义了基础常量,但你可以给键相关参数加key.前缀、值相关参数加value.前缀,以此实现分别指定不同的注册表和Schema名称,无需依赖单一配置项。


1. MSK Connect MySQL CDC连接器配置(生产者端)

在CDC连接器的配置中,为键和值的转换器分别指定独立的GSR参数:

# 键转换器配置
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=aws-glue-schema-registry://
key.converter.aws.schema.registry.registry.name=your-key-registry
key.converter.aws.schema.registry.schema.name=serverName.schemaName.tableName_key
key.converter.aws.schema.registry.schema.auto.register=true

# 值转换器配置
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=aws-glue-schema-registry://
value.converter.aws.schema.registry.registry.name=your-value-registry
value.converter.aws.schema.registry.schema.name=serverName.schemaName.tableName_value
value.converter.aws.schema.registry.schema.auto.register=true

2. Java生产者配置示例

针对键和值的序列化器,用前缀区分GSR配置:

import software.amazon.msk.connect.schemaregistry.AWSSchemaRegistryConstants;
import software.amazon.msk.connect.schemaregistry.serializers.avro.AwsAvroSerializer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;

import java.util.Properties;

public class GsrAvroProducer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-broker-endpoint");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, AwsAvroSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, AwsAvroSerializer.class.getName());

        // 键序列化器的GSR配置
        props.put("key." + AWSSchemaRegistryConstants.REGISTRY_NAME, "your-key-registry");
        props.put("key." + AWSSchemaRegistryConstants.SCHEMA_NAME, "serverName.schemaName.tableName_key");
        props.put("key." + AWSSchemaRegistryConstants.SCHEMA_AUTO_REGISTRATION_SETTING, "true");

        // 值序列化器的GSR配置
        props.put("value." + AWSSchemaRegistryConstants.REGISTRY_NAME, "your-value-registry");
        props.put("value." + AWSSchemaRegistryConstants.SCHEMA_NAME, "serverName.schemaName.tableName_value");
        props.put("value." + AWSSchemaRegistryConstants.SCHEMA_AUTO_REGISTRATION_SETTING, "true");

        // 发送AVRO消息(假设UserKey、UserValue是AVRO生成的实体类)
        try (KafkaProducer<Object, Object> producer = new KafkaProducer<>(props)) {
            ProducerRecord<Object, Object> record = new ProducerRecord<>("your-topic", new UserKey(), new UserValue());
            producer.send(record).get();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

3. Java消费者配置示例

反序列化器同样通过前缀区分键值的GSR配置:

import software.amazon.msk.connect.schemaregistry.AWSSchemaRegistryConstants;
import software.amazon.msk.connect.schemaregistry.deserializers.avro.AwsAvroDeserializer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class GsrAvroConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-broker-endpoint");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "your-consumer-group-id");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, AwsAvroDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, AwsAvroDeserializer.class.getName());

        // 键反序列化器的GSR配置
        props.put("key." + AWSSchemaRegistryConstants.REGISTRY_NAME, "your-key-registry");
        props.put("key." + AWSSchemaRegistryConstants.SCHEMA_NAME, "serverName.schemaName.tableName_key");

        // 值反序列化器的GSR配置
        props.put("value." + AWSSchemaRegistryConstants.REGISTRY_NAME, "your-value-registry");
        props.put("value." + AWSSchemaRegistryConstants.SCHEMA_NAME, "serverName.schemaName.tableName_value");

        // 拉取并处理消息
        try (KafkaConsumer<Object, Object> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList("your-topic"));
            while (true) {
                ConsumerRecords<Object, Object> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> {
                    UserKey key = (UserKey) record.key();
                    UserValue value = (UserValue) record.value();
                    // 业务逻辑处理
                    System.out.printf("Key: %s, Value: %s%n", key, value);
                });
            }
        }
    }
}

关键注意事项

  • 所有GSR配置参数(如SCHEMA_AUTO_REGISTRATION_SETTING、REGION等)都支持通过key./value.前缀实现键值分离配置。
  • 确保AVRO实体类的全限定名与你指定的SCHEMA_NAME完全匹配,否则会出现Schema匹配失败的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 07:10:28