使用camel-minio-sink-kafka-connector写入Avro消息到Minio遇类型转换失败
解决camel-minio-sink-kafka-connector无法转换Avro Struct到InputStream的问题
一、排查自定义类型转换器的注册环节
你写了转换器但没生效,大概率是注册步骤没做对,按以下步骤核对:
- 转换器类编写规范:必须用
@Converter注解标记类和转换方法,方法要明确接收org.apache.kafka.connect.data.Struct,返回java.io.InputStream。示例代码:
import org.apache.camel.Converter; import org.apache.kafka.connect.data.Struct; import org.apache.avro.io.Encoder; import org.apache.avro.io.EncoderFactory; import org.apache.avro.generic.GenericDatumWriter; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.IOException; @Converter public class StructToInputStreamConverter { @Converter public static InputStream convert(Struct struct) throws IOException { ByteArrayOutputStream baos = new ByteArrayOutputStream(); GenericDatumWriter<Struct> writer = new GenericDatumWriter<>(struct.schema()); Encoder encoder = EncoderFactory.get().binaryEncoder(baos, null); writer.write(struct, encoder); encoder.flush(); return new ByteArrayInputStream(baos.toByteArray()); } }
- 注册文件配置:在转换器jar包内创建
META-INF/services/org/apache/camel/type-converters文件,文件内一行写入转换器的全类名(比如com.yourpackage.StructToInputStreamConverter) - 类路径放置:将打包好的jar放到Kafka Connect的
plugins目录,或者camel-minio连接器的lib目录,确保连接器能加载到该类 - 连接器配置检查:确认连接器配置中
camel.component.minio.type-converter-enabled设为true(默认是true,但显式配置可避免意外)
二、更简单的替代方案:用Kafka Connect Avro转换器预处理
既然消息是Avro格式,直接让Kafka Connect先把Struct转成字节数组,Camel就能自动转成InputStream:
在连接器配置里添加以下参数:
value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://你的Schema Registry地址 value.converter.schemas.enable=false
配置后,Kafka Connect会把Avro Struct直接序列化为字节数组传给Minio连接器,Camel内置转换器可自动将字节数组转成InputStream,无需自定义转换器。
三、通过Camel路由模板注入处理逻辑
如果不想修改转换器,也可以在连接器路由里显式添加转换步骤:
在连接器配置中设置camel.sink.route-template:
camel.sink.route-template=from("{{camel.sink.path}}").process("structToInputStreamProcessor").to("{{camel.sink.to}}")
然后编写一个自定义Processor类(逻辑和转换器一致),打包后放到类路径,Camel会自动识别该Processor bean。
内容的提问来源于stack exchange,提问作者Andreas Baus
相关产品推荐
相关产品推荐

