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