Camel Salesforce Kafka Source Connector消息格式转换异常问题
问题根因
- NullPointerException异常触发原因:配置项
value.converter.schemas.enable设为true,但Camel Salesforce Source Connector输出的原始事件未附带Kafka Connect要求的Schema元数据结构,JsonConverter执行序列化逻辑时无法读取对应Schema信息,直接触发空指针异常。 - StringConverter输出非标准JSON的原因:默认配置下连接器会将Salesforce返回的JSON反序列化为Java Map对象,StringConverter会直接调用Map的toString()方法执行序列化,最终输出
key=value格式的非标准JSON串。
修复方案
调整连接器的以下配置项即可解决问题:
- 关闭JsonConverter的Schema启用开关,无需Kafka Connect封装Schema结构:
"value.converter.schemas.enable": "false" - 保留
camel.source.endpoint.rawPayload = true配置,确保连接器直接返回Salesforce输出的原生JSON字符串,不做Java对象转换。
若使用StringConverter作为值序列化器,完成上述配置后输出的原生JSON串可直接被下游Faust Agent正常解析;若使用JsonConverter作为值序列化器,关闭Schema开关后也可正常输出标准JSON格式数据。
修改后的完整配置参考如下:
{ "name": "sf_order_p_event_connector", "config": { "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "connector.class": "org.apache.camel.kafkaconnector.salesforce.CamelSalesforceSourceConnector", "camel.component.salesforce.loginUrl": "<redacted>", "camel.component.salesforce.instanceUrl": "<redacted>", "topics": "<redacted>", "camel.source.endpoint.rawPayload": "true", "camel.source.path.topicName": "/event/Order_Completed__e", "camel.source.endpoint.replayId": "-1", "camel.component.salesforce.authenticationType": "USERNAME_PASSWORD", "camel.component.salesforce.clientId": "<redacted>", "camel.component.salesforce.clientSecret": "<redacted>", "camel.component.salesforce.password": "<redacted>", "camel.component.salesforce.userName": "<redacted>", "camel.source.endpoint.apiVersion": "52.0" } }
内容的提问来源于stack exchange,提问作者user3554599
相关产品推荐
相关产品推荐

