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地址 } }
验证步骤
- 先启动Ignite服务端集群,确认所有节点正常上线
- 提交连接器配置,在控制中心确认连接器状态为
RUNNING - 往Kafka目标topic写入测试数据,查询Ignite对应缓存确认数据写入正常
内容的提问来源于stack exchange,提问作者Shubham Shirur
相关产品推荐
相关产品推荐

