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

Kafka-Ignite Sink Connector远程连接需要哪些配置?

Kafka Ignite Sink Connector 跨节点同步配置方案

前置检查

确认连接器所在节点和Ignite服务端节点网络互通,开放Ignite默认使用的端口段:

  • 发现端口:47500~47600
  • 通信端口:47100~47200

Ignite服务端节点Spring XML配置

核心配置项如下,示例见代码块:

  • 发现SPI绑定服务端对外可访问的IP,禁止绑定127.0.0.1
  • 使用静态IP发现器(TcpDiscoveryVmIpFinder),填入集群所有服务端节点的IP+端口,不要使用多播发现
  • 提前创建目标同步缓存,配置缓存名、键值类型和Kafka侧数据格式对齐
  • 若开启集群认证,需配置对应客户端接入的账号权限
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
       xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
    <bean id="ignite.cfg" class="org.apache.ignite.configuration.IgniteConfiguration">
        <!-- 发现SPI配置 -->
        <property name="discoverySpi">
            <bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
                <property name="localAddress" value="192.168.1.100"/> <!-- 替换为当前服务端对外IP -->
                <property name="localPort" value="47500"/>
                <property name="ipFinder">
                    <bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
                        <property name="addresses">
                            <list>
                                <!-- 替换为所有Ignite服务端节点的IP+端口 -->
                                <value>192.168.1.100:47500</value>
                                <value>192.168.1.101:47500</value>
                            </list>
                        </property>
                    </bean>
                </property>
            </bean>
        </property>

        <!-- 目标同步缓存配置 -->
        <property name="cacheConfiguration">
            <list>
                <bean class="org.apache.ignite.configuration.CacheConfiguration">
                    <property name="name" value="user_cache"/> <!-- 替换为你的缓存名 -->
                    <property name="keyType" value="java.lang.String"/>
                    <property name="valueType" value="com.example.UserPOJO"/> <!-- 替换为你的值类型 -->
                </bean>
            </list>
        </property>
    </bean>
</beans>

连接器侧传入的Spring XML配置

该XML是连接器作为Ignite厚客户端接入集群的配置,核心配置项如下:

  • 必须开启客户端模式,避免连接器作为数据节点加入Ignite集群
  • 发现SPI配置和服务端完全对齐,静态IP发现器的地址列表和服务端一致
  • 若服务端开启认证,需配置对应账号的凭证
  • 可开启对等类加载,避免自定义POJO的序列化类找不到问题
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
       xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
    <bean id="ignite.cfg" class="org.apache.ignite.configuration.IgniteConfiguration">
        <!-- 开启客户端模式 -->
        <property name="clientMode" value="true"/>
        <!-- 开启对等类加载,可选,自定义POJO场景建议开启 -->
        <property name="peerClassLoadingEnabled" value="true"/>

        <property name="discoverySpi">
            <bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
                <property name="localAddress" value="192.168.1.200"/> <!-- 替换为连接器所在节点的对外IP -->
                <property name="ipFinder">
                    <bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
                        <property name="addresses">
                            <list>
                                <!-- 和Ignite服务端的发现地址列表完全一致 -->
                                <value>192.168.1.100:47500</value>
                                <value>192.168.1.101:47500</value>
                            </list>
                        </property>
                    </bean>
                </property>
            </bean>
        </property>
    </bean>
</beans>

连接器核心属性配置

除了传入上述XML的路径外,核心配置示例如下:

{
  "name": "ignite-sink-connector",
  "config": {
    "connector.class": "org.apache.ignite.stream.kafka.connect.IgniteSinkConnector",
    "tasks.max": "3", // 最大不超过待同步Kafka topic的分区数
    "topics": "user_sync_topic", // 替换为你的Kafka topic名
    "spring.ignite.config.path": "/opt/kafka-connect/config/ignite-client-config.xml", // 替换为连接器节点上的XML绝对路径
    "ignite.cache.name": "user_cache", // 和服务端提前创建的缓存名完全一致
    "key.converter": "org.apache.kafka.connect.storage.StringConverter", // 和Kafka topic的key序列化格式对齐
    "value.converter": "io.confluent.connect.avro.AvroConverter", // 和Kafka topic的value序列化格式对齐
    "value.converter.schema.registry.url": "http://schema-registry:8081" // 若用Avro需配置schema registry地址
  }
}

验证步骤

  1. 先启动Ignite服务端集群,确认所有节点正常上线
  2. 提交连接器配置,在控制中心确认连接器状态为RUNNING
  3. 往Kafka目标topic写入测试数据,查询Ignite对应缓存确认数据写入正常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 12:45:04