[
https://issues.apache.org/jira/browse/IGNITE-16586?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Nikolay Izhikov updated IGNITE-16586:
-------------------------------------
Description:
Currently, only indexed parameters value can be provided for Cdc streamers.
We should support named parameters.
{code}
<bean id="cdc.streamer"
class="org.apache.ignite.cdc.IgniteToIgniteCdcStreamer">
<constructor-arg index="0">
<bean class="org.apache.ignite.configuration.IgniteConfiguration">
<property name="igniteInstanceName"
value="ignite-2029-cdc-client" />
<property name="clientMode" value="true" />
<property name="peerClassLoadingEnabled" value="true" />
<property name="localHost" value="127.0.0.1" />
<property name="discoverySpi">
<bean
class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
<property name="localPort" value="47600" />
<property name="ipFinder">
<bean
class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
<property name="addresses"
value="127.0.0.1:47600..47610" />
</bean>
</property>
<property name="joinTimeout" value="10000" />
</bean>
</property>
</bean>
</constructor-arg>
<constructor-arg index="1" value="false" />
<constructor-arg index="2">
<util:list>
<bean class="java.lang.String">
<constructor-arg type="String" value="terminator" />
</bean>
</util:list>
</constructor-arg>
<constructor-arg index="3" value="256" />
</bean>
{code}
was:
Now KafkaToIgniteCdcStreamerApplier[1] and IgniteToKafkaCdcStreamer[2] perform
requests with a hard-coded timeout equal to {{DFLT_REQ_TIMEOUT}}:
{code:title=KafkaToIgniteCdcStreamerApplier}
/** */
public static final int DFLT_REQ_TIMEOUT = 3;
...
private void poll(KafkaConsumer<Integer, byte[]> cnsmr) throws
IgniteCheckedException {
ConsumerRecords<Integer, byte[]> recs =
cnsmr.poll(Duration.ofSeconds(DFLT_REQ_TIMEOUT));
if (log.isDebugEnabled()) {
log.debug(
"Polled from consumer [assignments=" + cnsmr.assignment() +
",rcvdEvts=" + rcvdEvts.addAndGet(recs.count()) + ']'
);
}
apply(F.iterator(recs, this::deserialize, true, rec ->
F.isEmpty(caches) || caches.contains(rec.key())));
cnsmr.commitSync(Duration.ofSeconds(DFLT_REQ_TIMEOUT));
}
{code}
{code:title=IgniteToKafkaCdcStreamer}
/** Default kafka request timeout in seconds. */
public static final int DFLT_REQ_TIMEOUT = 5;
...
@Override public boolean onEvents(Iterator<CdcEvent> evts) {
List<Future<RecordMetadata>> futs = new ArrayList<>();
...
if (!futs.isEmpty()) {
try {
for (Future<RecordMetadata> fut : futs)
fut.get(DFLT_REQ_TIMEOUT, TimeUnit.SECONDS);
msgsSnt.add(futs.size());
lastMsgTs.value(System.currentTimeMillis());
}
{code}
We should have configurable timeout for requests to the Kafka.
#
https://github.com/apache/ignite-extensions/blob/master/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/kafka/KafkaToIgniteCdcStreamerApplier.java#L203
#
https://github.com/apache/ignite-extensions/blob/master/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/kafka/IgniteToKafkaCdcStreamer.java#L197
> Provide named parameters for Cdc streamers
> ------------------------------------------
>
> Key: IGNITE-16586
> URL: https://issues.apache.org/jira/browse/IGNITE-16586
> Project: Ignite
> Issue Type: Improvement
> Reporter: Nikolay Izhikov
> Priority: Minor
> Labels: IEP-59, ise
>
> Currently, only indexed parameters value can be provided for Cdc streamers.
> We should support named parameters.
> {code}
> <bean id="cdc.streamer"
> class="org.apache.ignite.cdc.IgniteToIgniteCdcStreamer">
> <constructor-arg index="0">
> <bean class="org.apache.ignite.configuration.IgniteConfiguration">
> <property name="igniteInstanceName"
> value="ignite-2029-cdc-client" />
> <property name="clientMode" value="true" />
> <property name="peerClassLoadingEnabled" value="true" />
> <property name="localHost" value="127.0.0.1" />
> <property name="discoverySpi">
> <bean
> class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
> <property name="localPort" value="47600" />
> <property name="ipFinder">
> <bean
> class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
> <property name="addresses"
> value="127.0.0.1:47600..47610" />
> </bean>
> </property>
> <property name="joinTimeout" value="10000" />
> </bean>
> </property>
> </bean>
> </constructor-arg>
> <constructor-arg index="1" value="false" />
> <constructor-arg index="2">
> <util:list>
> <bean class="java.lang.String">
> <constructor-arg type="String" value="terminator" />
> </bean>
> </util:list>
> </constructor-arg>
> <constructor-arg index="3" value="256" />
> </bean>
> {code}
--
This message was sent by Atlassian Jira
(v8.20.1#820001)