[ 
https://issues.apache.org/jira/browse/FLINK-23850?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17407235#comment-17407235
 ] 

Matthias commented on FLINK-23850:
----------------------------------

Thanks for the explaination, [~twalthr]. I will give it another try. Here are 
my findings of the initial test:

Setup
{code:bash}
# download necessary artifacts
wget https://archive.apache.org/dist/kafka/2.4.1/kafka_2.12-2.4.1.tgz
tar -xf kafka_2.12-2.4.1.tgz
wget 
https://dist.apache.org/repos/dist/dev/flink/flink-1.14.0-rc0/flink-1.14.0-bin-scala_2.12.tgz
tar -xf flink-1.14.0-bin-scala_2.12.tgz
wget 
https://repository.apache.org/content/repositories/orgapacheflink-1448/org/apache/flink/flink-sql-connector-kafka_2.12/1.14.0/flink-sql-connector-kafka_2.12-1.14.0.jar

# setup Kafka locally with two topics
./kafka_2.12-2.4.1/bin/zookeeper-server-start.sh 
kafka_2.12-2.4.1/config/zookeeper.properties
./kafka_2.12-2.4.1/bin/kafka-server-start.sh 
./kafka_2.12-2.4.1/config/server.properties

./kafka_2.12-2.4.1/bin/kafka-topics.sh --create --topic flink-23850-in 
--bootstrap-server localhost:9092
./kafka_2.12-2.4.1/bin/kafka-topics.sh --create --topic flink-23850-out 
--bootstrap-server localhost:9092

# create consumer on both topics for verification purposes
./kafka_2.12-2.4.1/bin/kafka-console-consumer.sh --topic flink-23850-in 
--from-beginning --bootstrap-server localhost:9092
./kafka_2.12-2.4.1/bin/kafka-console-consumer.sh --topic flink-23850-out 
--from-beginning --bootstrap-server localhost:9092

# create producer for manual content generation
./kafka_2.12-2.4.1/bin/kafka-console-producer.sh --topic flink-23850-in 
--broker-list localhost:9092

# start Flink cluster with two TaskManagers
./flink-1.14.0/bin/start-cluster.sh
./flink-1.14.0/bin/taskmanager.sh start

# start SQLClient with SQL connector for Kafka
./flink-1.14.0/bin/sql-client.sh --jar flink-sql-connector-kafka_2.12-1.14.0.jar
{code}

Table API Queries tested:
{code:sql}
-- create input table
CREATE TABLE FLINK_23850_IN (
  text STRING, 
  proc_time AS PROCTIME()
) WITH (
  'connector' = 'kafka', 
  'topic' = 'flink-23850-in', 
  'properties.bootstrap.servers' = 'localhost:9092', 
  'format' = 'csv', 
  'scan.startup.mode' = 'earliest-offset',
  'properties.client.id' = 'flink-23850-input'
);

-- create output table
CREATE TABLE FLINK_23850_OUT (
  text STRING
) WITH (
  'connector' = 'kafka', 
  'topic' = 'flink-23850-out', 
  'properties.bootstrap.servers' = 'localhost:9092', 
  'format' = 'csv', 
  'scan.startup.mode' = 'earliest-offset',
  'properties.client.id' = 'flink-23850-output'
);

-- creates job for forwarding the events
INSERT INTO FLINK_23850_OUT SELECT text FROM FLINK_23850_IN;

-- creates monitoring query for input
SELECT * FROM FLINK_23850_IN;
{code}

The logs don't indicate anything suspicious. {{KafkaSource}}/{{KafkaSink}} seem 
to be used based on the logs:
{code}
2021-08-31 12:34:36,526 WARN  
org.apache.flink.connector.kafka.sink.KafkaSinkBuilder       [] - Property 
[transaction.timeout.ms] not specified. Setting it to PT1H
[...]
2021-08-31 12:35:36,284 WARN  
org.apache.flink.streaming.connectors.kafka.table.KafkaDynamicSource [] - 
Property "group.id" is required for offset commit but not set in table options. 
Assigning "KafkaSource-default_catalog.default_database.FLINK_23850_IN" as 
consumer group id
{code}

I ran into an issue, though, when cancelling the {{SELECT}} job:
{{code}}
2021-08-31 12:34:37,091 WARN  
org.apache.flink.streaming.api.operators.collect.CollectResultFetcher [] - An 
exception occurred when fetching query results
java.util.concurrent.ExecutionException: 
org.apache.flink.runtime.rest.util.RestClientException: [Internal server 
error., <Exception on server side:
org.apache.flink.runtime.messages.FlinkJobNotFoundException: Could not find 
Flink job (3df388d05273cc9782458228081ea5bf)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.getJobMasterGateway(Dispatcher.java:909)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.performOperationOnJobMasterGateway(Dispatcher.java:923)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.deliverCoordinationRequestToCoordinator(Dispatcher.java:719)
        at sun.reflect.GeneratedMethodAccessor24.invoke(Unknown Source)
        at 
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.lambda$handleRpcInvocation$1(AkkaRpcActor.java:316)
        at 
org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:83)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:314)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:217)
        at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:78)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:163)
        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24)
        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20)
        at scala.PartialFunction.applyOrElse(PartialFunction.scala:123)
        at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122)
        at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
        at akka.actor.Actor.aroundReceive(Actor.scala:537)
        at akka.actor.Actor.aroundReceive$(Actor.scala:535)
        at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:220)
        at akka.actor.ActorCell.receiveMessage(ActorCell.scala:580)
        at akka.actor.ActorCell.invoke(ActorCell.scala:548)
        at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:270)
        at akka.dispatch.Mailbox.run(Mailbox.scala:231)
        at akka.dispatch.Mailbox.exec(Mailbox.scala:243)
        at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289)
        at 
java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056)
        at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692)
        at 
java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:175)

End of exception on server side>]
        at 
java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357) 
~[?:1.8.0_265]
        at 
java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908) 
~[?:1.8.0_265]
        at 
org.apache.flink.streaming.api.operators.collect.CollectResultFetcher.sendRequest(CollectResultFetcher.java:163)
 ~[flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.streaming.api.operators.collect.CollectResultFetcher.next(CollectResultFetcher.java:128)
 [flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.streaming.api.operators.collect.CollectResultIterator.nextResultFromFetcher(CollectResultIterator.java:106)
 [flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.streaming.api.operators.collect.CollectResultIterator.hasNext(CollectResultIterator.java:80)
 [flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.table.api.internal.TableResultImpl$CloseableRowIteratorWrapper.hasNext(TableResultImpl.java:370)
 [flink-table_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.table.client.gateway.local.result.CollectResultBase$ResultRetrievalThread.run(CollectResultBase.java:74)
 [flink-sql-client_2.12-1.14.0.jar:1.14.0]
Caused by: org.apache.flink.runtime.rest.util.RestClientException: [Internal 
server error., <Exception on server side:
org.apache.flink.runtime.messages.FlinkJobNotFoundException: Could not find 
Flink job (3df388d05273cc9782458228081ea5bf)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.getJobMasterGateway(Dispatcher.java:909)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.performOperationOnJobMasterGateway(Dispatcher.java:923)
        at 
org.apache.flink.runtime.dispatcher.Dispatcher.deliverCoordinationRequestToCoordinator(Dispatcher.java:719)
        at sun.reflect.GeneratedMethodAccessor24.invoke(Unknown Source)
        at 
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.lambda$handleRpcInvocation$1(AkkaRpcActor.java:316)
        at 
org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:83)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:314)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:217)
        at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:78)
        at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:163)
        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24)
        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20)
        at scala.PartialFunction.applyOrElse(PartialFunction.scala:123)
        at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122)
        at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
        at akka.actor.Actor.aroundReceive(Actor.scala:537)
        at akka.actor.Actor.aroundReceive$(Actor.scala:535)
        at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:220)
        at akka.actor.ActorCell.receiveMessage(ActorCell.scala:580)
        at akka.actor.ActorCell.invoke(ActorCell.scala:548)
        at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:270)
        at akka.dispatch.Mailbox.run(Mailbox.scala:231)
        at akka.dispatch.Mailbox.exec(Mailbox.scala:243)
        at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289)
        at 
java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056)
        at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692)
        at 
java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:175)

End of exception on server side>]
        at 
org.apache.flink.runtime.rest.RestClient.parseResponse(RestClient.java:532) 
~[flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
org.apache.flink.runtime.rest.RestClient.lambda$submitRequest$3(RestClient.java:512)
 ~[flink-dist_2.12-1.14.0.jar:1.14.0]
        at 
java.util.concurrent.CompletableFuture.uniCompose(CompletableFuture.java:966) 
~[?:1.8.0_265]
        at 
java.util.concurrent.CompletableFuture$UniCompose.tryFire(CompletableFuture.java:940)
 ~[?:1.8.0_265]
        at 
java.util.concurrent.CompletableFuture$Completion.run(CompletableFuture.java:456)
 ~[?:1.8.0_265]
        at 
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) 
~[?:1.8.0_265]
        at 
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) 
~[?:1.8.0_265]
        at java.lang.Thread.run(Thread.java:748) ~[?:1.8.0_265]
{{code}}

Is this a known issue? I cancelled the query with typing {{Q}} as requested in 
the CLI?

> Test Kafka table connector with new runtime provider
> ----------------------------------------------------
>
>                 Key: FLINK-23850
>                 URL: https://issues.apache.org/jira/browse/FLINK-23850
>             Project: Flink
>          Issue Type: Sub-task
>          Components: Connectors / Kafka
>    Affects Versions: 1.14.0
>            Reporter: Qingsheng Ren
>            Assignee: Matthias
>            Priority: Blocker
>              Labels: release-testing
>             Fix For: 1.14.0
>
>
> The runtime provider of Kafka table connector has been replaced with new 
> KafkaSource and KafkaSink. The table connector requires to be tested to make 
> sure nothing is surprised to Table/SQL API users. 



--
This message was sent by Atlassian Jira
(v8.3.4#803005)

Reply via email to