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

WSO2 Micro Integrator连接Kafka时出现Kafka连接器错误

解决WSO2 MI Kafka连接器"Connection name is not set"错误

我按照官方文档尝试将Kafka与WSO2 Micro Integrator(MI 4.1.0)连接,使用WSO2 Integration Studio开发了如下API代码,使用的Kafka版本为2.12,Kafka连接器版本为3.12:

<?xml version="1.0" encoding="UTF-8"?>
<api context="/create-customer" name="create-customer" xmlns="http://ws.apache.org/ns/synapse">
    <resource methods="POST">
        <inSequence>
            <kafkaTransport.init>
                <bootstrapServers>localhost:9092</bootstrapServers>
                <keySerializerClass>org.apache.kafka.common.serialization.StringSerializer</keySerializerClass>
                <valueSerializerClass>org.apache.kafka.common.serialization.StringSerializer</valueSerializerClass>
            </kafkaTransport.init>
            <kafkaTransport.publishMessages>
                <topic>customer</topic>
            </kafkaTransport.publishMessages>
            <respond/>
        </inSequence>
        <outSequence/>
        <faultSequence/>
    </resource>
</api>

发送请求时出现以下错误:

[2023-01-14 21:27:36,570]  INFO {KafkaProduceConnector} - {api:create-customer} SEND : send message to  Broker lists
[2023-01-14 21:27:36,583] ERROR {KafkaProduceConnector} - {api:create-customer} Kafka producer connector : Error sending the message to broker 
org.wso2.carbon.connector.exception.InvalidConfigurationException: Connection name is not set.
    at org.wso2.carbon.connector.KafkaProduceConnector.getConnectionName(KafkaProduceConnector.java:262)
    at org.wso2.carbon.connector.KafkaProduceConnector.publishMessage(KafkaProduceConnector.java:237)
    at org.wso2.carbon.connector.KafkaProduceConnector.connect(KafkaProduceConnector.java:138)
    at org.wso2.carbon.connector.core.AbstractConnector.mediate(AbstractConnector.java:32)
    at org.apache.synapse.mediators.ext.ClassMediator.updateInstancePropertiesAndMediate(ClassMediator.java:178)
    at org.apache.synapse.mediators.ext.ClassMediator.mediate(ClassMediator.java:97)
    at org.apache.synapse.mediators.AbstractListMediator.mediate(AbstractListMediator.java:110)
    at org.apache.synapse.mediators.AbstractListMediator.mediate(AbstractListMediator.java:72)
    at org.apache.synapse.mediators.template.TemplateMediator.mediate(TemplateMediator.java:136)
    at org.apache.synapse.mediators.template.InvokeMediator.mediate(InvokeMediator.java:170)
    at org.apache.synapse.mediators.template.InvokeMediator.mediate(InvokeMediator.java:93)
    at org.apache.synapse.mediators.AbstractListMediator.mediate(AbstractListMediator.java:110)
    at org.apache.synapse.mediators.AbstractListMediator.mediate(AbstractListMediator.java:72)
    at org.apache.synapse.mediators.base.SequenceMediator.mediate(SequenceMediator.java:158)
    at org.apache.synapse.api.Resource.process(Resource.java:342)
    at org.apache.synapse.api.API.process(API.java:477)
    at org.apache.synapse.api.AbstractApiHandler.apiProcess(AbstractApiHandler.java:93)
    at org.apache.synapse.api.AbstractApiHandler.dispatchToAPI(AbstractApiHandler.java:71)
    at org.apache.synapse.api.rest.RestRequestHandler.dispatchToAPI(RestRequestHandler.java:90)
    at org.apache.synapse.api.rest.RestRequestHandler.process(RestRequestHandler.java:76)
    at org.apache.synapse.rest.RESTRequestHandler.process(RESTRequestHandler.java:54)
    at org.apache.synapse.core.axis2.Axis2SynapseEnvironment.injectMessage(Axis2SynapseEnvironment.java:344)
    at org.apache.synapse.core.axis2.SynapseMessageReceiver.receive(SynapseMessageReceiver.java:101)
    at org.apache.axis2.engine.AxisEngine.receive(AxisEngine.java:180)
    at org.apache.synapse.transport.passthru.ServerWorker.processNonEntityEnclosingRESTHandler(ServerWorker.java:376)
    at org.apache.synapse.transport.passthru.ServerWorker.processEntityEnclosingRequest(ServerWorker.java:435)
    at org.apache.synapse.transport.passthru.ServerWorker.run(ServerWorker.java:183)
    at org.apache.axis2.transport.base.threads.NativeWorkerPool$1.run(NativeWorkerPool.java:172)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at java.base/java.lang.Thread.run(Thread.java:829)

根据配置文档,<kafkaTransport.init>中可设置连接名,但Integration Studio的该组件属性面板里没有这个字段,请问如何解决?


解决方案

  • 直接手动编辑XML代码,在<kafkaTransport.init>节点中添加connectionName配置项,无需依赖可视化面板。
  • 同时在<kafkaTransport.publishMessages>节点中指定相同的connectionName,确保生产者能关联到初始化的连接。

修改后的完整配置:

<?xml version="1.0" encoding="UTF-8"?>
<api context="/create-customer" name="create-customer" xmlns="http://ws.apache.org/ns/synapse">
    <resource methods="POST">
        <inSequence>
            <kafkaTransport.init>
                <connectionName>KafkaProducerConn</connectionName>
                <bootstrapServers>localhost:9092</bootstrapServers>
                <keySerializerClass>org.apache.kafka.common.serialization.StringSerializer</keySerializerClass>
                <valueSerializerClass>org.apache.kafka.common.serialization.StringSerializer</valueSerializerClass>
            </kafkaTransport.init>
            <kafkaTransport.publishMessages>
                <topic>customer</topic>
                <connectionName>KafkaProducerConn</connectionName>
            </kafkaTransport.publishMessages>
            <respond/>
        </inSequence>
        <outSequence/>
        <faultSequence/>
    </resource>
</api>
  • 说明:Integration Studio的可视化面板有时不会展示所有连接器属性,直接编辑XML是WSO2开发中常见的操作,添加的配置项会被MI正常识别并加载。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 02:35:48