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
相关产品推荐
相关产品推荐

