[ 
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)

Reply via email to