Flink CEP????????????????flink??????????????cep????

2021-04-15 Thread Asahi Lee
hi??
      flink 
cep??cep

flink 1.13.0 ??????flink sql ??????????????????????????????????schema.name

2021-05-18 Thread Asahi Lee
hi!
      flink jdbc ?? 
table-name??
CREATE TABLE MyUserTable (   id BIGINT,   name STRING,   age INT,   status 
BOOLEAN,   PRIMARY KEY (id) NOT ENFORCED ) WITH ('connector' = 'jdbc',
'url' = 'jdbc:mysql://localhost:3306/mydatabase','table-name' = 'org.users' 
);

????flink 1.12 ????????Flink SQL ???? ??sink?? ??????????????

2021-08-25 Thread ????
hi all: 
    ??flink 1.12 sql  ?? 
  kafka ??msyql ?? 
kafkadb??kafka ??
   ??
   insert into db_table_sink select * from  
kafka_source_table;
   insert into kafka_table_sink select * from kafka_source_table;


  flink SQL ?? 
??db??kafka??flink ??

?????? ????flink 1.12 ????????Flink SQL ???? ??sink?? ??????????????

2021-08-25 Thread ????
statement set[1]  StatementSet.addInsertSql 
??sql execute()




--  --
??: 
   "user-zh"

https://ci.apache.org/projects/flink/flink-docs-master/docs/dev/table/sqlclient/#execute-a-set-of-sql-statements

 

?????? ????flink 1.12 ????????Flink SQL ???? ??sink?? ??????????????

2021-08-25 Thread ????
flink-connector-kafka_2.11 ?? flink-connector-jdbc_2.11?? 
??mysql ?? ?? ?? 
kafka??java.sql.BatchUpdateException 
??3 
sink Kafka ??Kafka??  'sink.semantic' = 
'exactly-once', consumer  --isolation-level read_committed 
sink db ??sink kafka??flink 
??




--  --
??: 
   "user-zh"

https://github.com/ververica/flink-cdc-connectors

 https://ci.apache.org/projects/flink/flink-docs-master/docs/dev/table/sqlclient/#execute-a-set-of-sql-statements
>
>  

flink???????????????? flink last checkpoint ???????????????????????????? state ??????????????

2020-06-22 Thread ????????
??flink?? 
??checkpoint  

?????? flink???????????????? flink last checkpoint ???????????????????????????? state ??????????????

2020-06-23 Thread cs
??checkpoint??
StreamExecutionEnvironment.getCheckpointConfig().setFailOnCheckpointingErrors(true);
/**
 * Sets the expected behaviour for tasks in case that they encounter an error 
in their checkpointing procedure.
 * If this is set to true, the task will fail on checkpointing error. If this 
is set to false, the task will only
 * decline a the checkpoint and continue running. The default is true.
 */
public void setFailOnCheckpointingErrors(boolean failOnCheckpointingErrors) {
   this.failOnCheckpointingErrors = failOnCheckpointingErrors;
}


--  --
??: "LakeShen"

??????flink????????

2020-07-08 Thread Yichao Yang
Hi,


keyby?? warningPojo + String 
??keybykey??

 timestamp assigner 



Best,
Yichao Yang




--  --
??: "??"

??????flink????????

2020-07-09 Thread ????????
new ProcessWindowFunction?
???




--  --
??: 
   "user-zh"



????????????flink????????????????

2020-07-20 Thread jiafu
flinkflink-1.8.1
org.apache.flink.runtime.executiongraph.ExecutionGraphException: Trying to 
eagerly schedule a task whose inputs are not ready (result type: 
PIPELINED_BOUNDED, partition consumable: false, producer state: SCHEDULED, 
producer slot: null).at 
org.apache.flink.runtime.deployment.InputChannelDeploymentDescriptor.fromEdges(InputChannelDeploymentDescriptor.java:145)
at 
org.apache.flink.runtime.executiongraph.ExecutionVertex.createDeploymentDescriptor(ExecutionVertex.java:840)
 at 
org.apache.flink.runtime.executiongraph.Execution.deploy(Execution.java:621)
 at 
org.apache.flink.util.function.ThrowingRunnable.lambda$unchecked$0(ThrowingRunnable.java:50)
 at 
java.util.concurrent.CompletableFuture.uniRun(CompletableFuture.java:705)at 
java.util.concurrent.CompletableFuture.uniRunStage(CompletableFuture.java:717)  
 at 
java.util.concurrent.CompletableFuture.thenRun(CompletableFuture.java:2010)  at 
org.apache.flink.runtime.executiongraph.Execution.scheduleForExecution(Execution.java:436)
   at 
org.apache.flink.runtime.executiongraph.ExecutionVertex.scheduleForExecution(ExecutionVertex.java:637)
   at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.restart(FailoverRegion.java:229)
 at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.reset(FailoverRegion.java:186)
   at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.allVerticesInTerminalState(FailoverRegion.java:96)
   at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.lambda$cancel$0(FailoverRegion.java:146)
 at 
java.util.concurrent.CompletableFuture.uniAccept(CompletableFuture.java:656)
 at 
java.util.concurrent.CompletableFuture$UniAccept.tryFire(CompletableFuture.java:632)
 at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) 
 at 
java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1962)
 at 
org.apache.flink.runtime.concurrent.FutureUtils$WaitingConjunctFuture.handleCompletedFuture(FutureUtils.java:633)
at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
   at 
java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:736)
   at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) 
 at 
java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1962)
 at 
org.apache.flink.runtime.executiongraph.Execution.lambda$releaseAssignedResource$11(Execution.java:1350)
 at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
   at 
java.util.concurrent.CompletableFuture.uniWhenCompleteStage(CompletableFuture.java:778)
  at 
java.util.concurrent.CompletableFuture.whenComplete(CompletableFuture.java:2140)
 at 
org.apache.flink.runtime.executiongraph.Execution.releaseAssignedResource(Execution.java:1345)
   at 
org.apache.flink.runtime.executiongraph.Execution.finishCancellation(Execution.java:1115)
at 
org.apache.flink.runtime.executiongraph.Execution.completeCancelling(Execution.java:1094)
at 
org.apache.flink.runtime.executiongraph.ExecutionGraph.updateState(ExecutionGraph.java:1628)
 at 
org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:517)
at sun.reflect.GeneratedMethodAccessor63.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.handleRpcInvocation(AkkaRpcActor.java:274)
at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:189)
   at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)
at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.onReceive(AkkaRpcActor.java:147) 
 at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.onReceive(FencedAkkaRpcActor.java:40)
   at 
akka.actor.UntypedActor$$anonfun$receive$1.applyOrElse(UntypedActor.scala:165)  
 at akka.actor.Actor$class.aroundReceive(Actor.scala:502)at 
akka.actor.UntypedActor.aroundReceive(UntypedActor.scala:95) at 
akka.actor.ActorCell.receiveMessage(ActorCell.scala:526) at 
akka.actor.ActorCell.invoke(ActorCell.scala:495) at 
akka.dispatch.Mailbox.processMailbox(Mailbox.scala:257)  at 
akka.dispatch.Mailbox.run(Mailbox.scala:224) at 
akka.dispatch.Mailbox.exec(Mailbox.scala:234)at 
scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at 
scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
 at 
scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at 
scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

????: ????????????flink????????????????

2020-07-20 Thread sjlsumait...@163.com
??



sjlsumait...@163.com
 
 jiafu
?? 2020-07-20 19:31
 user-zh
?? flink
flinkflink-1.8.1
org.apache.flink.runtime.executiongraph.ExecutionGraphException: Trying to 
eagerly schedule a task whose inputs are not ready (result type: 
PIPELINED_BOUNDED, partition consumable: false, producer state: SCHEDULED, 
producer slot: null). at 
org.apache.flink.runtime.deployment.InputChannelDeploymentDescriptor.fromEdges(InputChannelDeploymentDescriptor.java:145)
 at 
org.apache.flink.runtime.executiongraph.ExecutionVertex.createDeploymentDescriptor(ExecutionVertex.java:840)
 at 
org.apache.flink.runtime.executiongraph.Execution.deploy(Execution.java:621) at 
org.apache.flink.util.function.ThrowingRunnable.lambda$unchecked$0(ThrowingRunnable.java:50)
 at java.util.concurrent.CompletableFuture.uniRun(CompletableFuture.java:705) 
at 
java.util.concurrent.CompletableFuture.uniRunStage(CompletableFuture.java:717) 
at java.util.concurrent.CompletableFuture.thenRun(CompletableFuture.java:2010) 
at 
org.apache.flink.runtime.executiongraph.Execution.scheduleForExecution(Execution.java:436)
 at 
org.apache.flink.runtime.executiongraph.ExecutionVertex.scheduleForExecution(ExecutionVertex.java:637)
 at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.restart(FailoverRegion.java:229)
 at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.reset(FailoverRegion.java:186)
 at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.allVerticesInTerminalState(FailoverRegion.java:96)
 at 
org.apache.flink.runtime.executiongraph.failover.FailoverRegion.lambda$cancel$0(FailoverRegion.java:146)
 at 
java.util.concurrent.CompletableFuture.uniAccept(CompletableFuture.java:656) at 
java.util.concurrent.CompletableFuture$UniAccept.tryFire(CompletableFuture.java:632)
 at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) 
at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1962) 
at 
org.apache.flink.runtime.concurrent.FutureUtils$WaitingConjunctFuture.handleCompletedFuture(FutureUtils.java:633)
 at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
 at 
java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:736)
 at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) 
at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1962) 
at 
org.apache.flink.runtime.executiongraph.Execution.lambda$releaseAssignedResource$11(Execution.java:1350)
 at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
 at 
java.util.concurrent.CompletableFuture.uniWhenCompleteStage(CompletableFuture.java:778)
 at 
java.util.concurrent.CompletableFuture.whenComplete(CompletableFuture.java:2140)
 at 
org.apache.flink.runtime.executiongraph.Execution.releaseAssignedResource(Execution.java:1345)
 at 
org.apache.flink.runtime.executiongraph.Execution.finishCancellation(Execution.java:1115)
 at 
org.apache.flink.runtime.executiongraph.Execution.completeCancelling(Execution.java:1094)
 at 
org.apache.flink.runtime.executiongraph.ExecutionGraph.updateState(ExecutionGraph.java:1628)
 at 
org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:517)
 at sun.reflect.GeneratedMethodAccessor63.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.handleRpcInvocation(AkkaRpcActor.java:274)
 at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:189)
 at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)
 at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.onReceive(AkkaRpcActor.java:147) 
at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.onReceive(FencedAkkaRpcActor.java:40)
 at 
akka.actor.UntypedActor$$anonfun$receive$1.applyOrElse(UntypedActor.scala:165) 
at akka.actor.Actor$class.aroundReceive(Actor.scala:502) at 
akka.actor.UntypedActor.aroundReceive(UntypedActor.scala:95) at 
akka.actor.ActorCell.receiveMessage(ActorCell.scala:526) at 
akka.actor.ActorCell.invoke(ActorCell.scala:495) at 
akka.dispatch.Mailbox.processMailbox(Mailbox.scala:257) at 
akka.dispatch.Mailbox.run(Mailbox.scala:224) at 
akka.dispatch.Mailbox.exec(Mailbox.scala:234) at 
scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at 
scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
 at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at 
scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)


????????flink??????????????????

2020-08-05 Thread samuel....@ubtrobot.com
flink 
,??


   ??mysql??json
{"times":5}  ---5??
{"temperature": 80} ---80
   
  1)????kafka
  
2)flinkkafka??


??
1. ????
2.??flink CEP??
3.??


   

 


?????? ????????flink??????????????????

2020-08-05 Thread ????????
cep
cepGroupPattern??within??wait??
cepcep??cep
?? cepflink1.7?? 
 https://developer.aliyun.com/article/738451





--  --
??: 
   "user-zh"

https://blog.csdn.net/zhangjun5965/article/details/106573528

samuel@ubtrobot.com 

????????????????flink

2020-08-12 Thread ??????
kafka0.10??flink1.10.flinkkafka

??????????flink????

2020-08-13 Thread ??????(Jiacheng Jiang)
1.10job??--  --
??: ""

Flink??????????????????????

2020-08-25 Thread Sun_yijia
??A??B??AB??
??B??ABA


??Flink??AB

Flink??????????????????????

2020-08-25 Thread Sun_yijia



??A??B??AB??
??B??ABA


??FlinkAB

?????? Flink??????????????????????

2020-08-27 Thread Sun_yijia
~


--  --
??: 
   "user-zh"



??Flink??????????????

2020-08-29 Thread ????????
hi,all:


??demoflink?
.

??????????????flink??????????????????

2020-09-04 Thread ????
??


   Flink+drools drools
2020-9-4
| |

|
|
hold_li...@163.com
|
??
??2020??8??6?? 10:26??samuel@ubtrobot.com ??
flink 
,??


??mysql??json
{"times":5}  ---5??
{"temperature": 80} ---80

1)????kafka
2)flinkkafka??


??
1. ????
2.??flink CEP??
3.??







flink??????????????????

2020-10-28 Thread ????????????
??
     Linuxflink??31.9.3 
 ??1?? ./start-cluster.sh  Name or service not 
knownname ce-hjjcgl-svr-02??21??3??
 [appuser@ce-hjjcgl-svr-01 bin]$ ./start-cluster.sh
 Starting cluster.
 Starting standalonesession daemon on host ce-hjjcgl-svr-01.
 : Name or service not knownname ce-hjjcgl-svr-02
 Starting taskexecutor daemon on host ce-hjjcgl-svr-03.


 
 
 
 ??2??./start-cluster.sh??
 [appuser@ce-hjjcgl-svr-02 bin]$ ./start-cluster.sh
Starting cluster.
[INFO] 1 instance(s) of standalonesession are already running on 
ce-hjjcgl-svr-02.
Starting standalonesession daemon on host ce-hjjcgl-svr-02.
: Name or service not knownname ce-hjjcgl-svr-02
 


 
 ??sshhosts??
 127.0.0.1   localhost localhost.localdomain localhost4 
localhost4.localdomain4
::1         localhost localhost.localdomain 
localhost6 localhost6.localdomain6
10.212.139.219 ce-hjjcgl-svr-01
10.212.139.220 ce-hjjcgl-svr-02
10.212.139.221 ce-hjjcgl-svr-03

??masters??
ce-hjjcgl-svr-01:8081

slaves??
ce-hjjcgl-svr-02
ce-hjjcgl-svr-03


2??flink-conf.yaml??

env.java.home: /usr/java/jdk1.8.0_202/


jobmanager.rpc.address: ce-hjjcgl-svr-02
 
   jobmanager.rpc.port: 6123



??

flink????????????

2020-11-07 Thread ????(Bob Hu)
??flink??yarn 
killflink1.11.0??
bin/flink run -m yarn-cluster -yjm 2048m -ytm 8192m -ys 2 
xxx.jar,rocksdb??taskmanager.memory.managed.fraction=0.6;taskmanager.memory.jvm-overhead.fraction=0.2flink??taskmanage??jarnativenative??


Free Slots / All Slots:0 / 2
CPU Cores:24
Physical Memory:251 GB
JVM Heap Size:1.82 GB
Flink Managed Memory:4.05 GB

Memory


JVM (Heap/Non-Heap)


Type
Committed
Used
Maximum

Heap1.81 GB1.13 GB1.81 GB
Non-Heap169 MB160 MB1.48 GB
Total1.98 GB1.29 GB3.30 GB





Outside JVM


Type
Count
Used
Capacity

Direct24,493718 MB718 MB
Mapped00 B0 B






Network


Memory Segments


Type
Count

Available21,715
Total22,118





Garbage Collection


Collector
Count
Time

PS_Scavenge19917,433
PS_MarkSweep44,173

flink????????????

2021-01-05 Thread Waldeinsamkeit.


flink??????????

2021-02-21 Thread liujian
Hi:
    
??,??3??,??,,??,kudu??
    
    :
      
 jdbc-connecter,join,??,

??????flink??????????

2021-02-21 Thread liujian
??,,??,flinkjoin??,hbase,??




--  --
??: 
   "user-zh"



flink????????????????

2021-04-24 Thread Natasha
hi,all
??
tableEnv.createTemporaryView("sensor", sensorTable);
val resultSqlTable = tableEnv.sqlQuery("select country, count(order_id) as cnt from sensor group by country");

??socket??
001 usa
002 usa
003 china
002 china
004 usa

??
usa, 2
china, 2

??
usa, 3
china, 2

??usa??usa

??

thanks

Flink????????????????

2021-04-24 Thread Natasha
hi,all
??
tableEnv.createTemporaryView("sensor", sensorTable);
val resultSqlTable = tableEnv.sqlQuery("select country, count(order_id) as cnt 
from sensor group by country");


??socket??
001 usa
002 usa
003 china
002 china
004 usa


??
usa, 2
china, 2


??
usa, 3
china, 2


??usa??usa

??


thanks

Flink????????????????

2021-04-24 Thread Natasha
hi,all
??
tableEnv.createTemporaryView("sensor", sensorTable);
val resultSqlTable = tableEnv.sqlQuery("select country, count(order_id) as cnt 
from sensor group by country");
tableEnv.toRetractStream[WaterSensorCnt](resultSqlTable).print("result");


??socket??
001 usa
002 usa
003 china
002 china
004 usa


??
usa, 2
china, 2


??
usa, 3
china, 2


??usa??usa

??


thanks

??????Flink????????????????

2021-04-24 Thread ??????
??

??2021??04??25?? 11:50??Natasha ??
hi,all
??
tableEnv.createTemporaryView("sensor", sensorTable);
val resultSqlTable = tableEnv.sqlQuery("select country, count(order_id) as cnt 
from sensor group by country");
tableEnv.toRetractStream[WaterSensorCnt](resultSqlTable).print("result");


??socket??
001 usa
002 usa
003 china
002 china
004 usa


??
usa, 2
china, 2


??
usa, 3
china, 2


??usa??usa

??


thanks

??????Flink????????????????

2021-04-24 Thread lian



??2021??04??25?? 11:47??Natasha ??
hi,all
??
tableEnv.createTemporaryView("sensor", sensorTable);
val resultSqlTable = tableEnv.sqlQuery("select country, count(order_id) as cnt 
from sensor group by country");


??socket??
001 usa
002 usa
003 china
002 china
004 usa


??
usa, 2
china, 2


??
usa, 3
china, 2


??usa??usa

??


thanks

flink ????

2021-05-27 Thread liujian
flink 
lookup??,lookup,?


?????? flink ????

2021-05-27 Thread liujian
Hi:
??Leonard,??flink cdc??(??), join ??(debezium -> 
kafka), eventTime , ??eventTime  
lookupeventTime
1, changlog,  debezium -> kafka , 
??flink cdc, ????flink cdc??,????flink 
cdc??watermark?
2, ??watermark??? 


??


--  --
??: 
   "user-zh"

http://a.id/>; = B.id <http://b.id/>;  
??(proctime??)


,
Leonard

> ?? 2021??5??2716:35??liujian <13597820...@qq.com> ??
> 
> flink 
lookup??,lookup,?
> 

Flink????????-??????????????????????????????????????????

2021-06-10 Thread ??????
keyed ??trigger


key??24onElement()FIRE_AND_PURGEvaluestate??onEventTime()2??



1??keyby??watermark??onEventTime()keykey?? 
 
key??valuestate??


2onElement()onEventTime()TriggerResult??




??
??
class MyTrigger extends Trigger

flink ??????????????

2021-07-13 Thread ????????????????
Hi All??


    ??Flink 
checkpoint??2min??
??2min??  ??
   



 The program finished with the following exception:


org.apache.flink.util.FlinkException: Triggering a savepoint for the job 
 failed.
at 
org.apache.flink.client.cli.CliFrontend.triggerSavepoint(CliFrontend.java:777)
at 
org.apache.flink.client.cli.CliFrontend.lambda$savepoint$9(CliFrontend.java:754)
at 
org.apache.flink.client.cli.CliFrontend.runClusterAction(CliFrontend.java:1002)
at 
org.apache.flink.client.cli.CliFrontend.savepoint(CliFrontend.java:751)
at 
org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1072)
at 
org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1132)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at 
org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1730)
at 
org.apache.flink.runtime.security.contexts.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1132)
Caused by: java.util.concurrent.TimeoutException
at 
org.apache.flink.runtime.concurrent.FutureUtils$Timeout.run(FutureUtils.java:1255)
at 
org.apache.flink.runtime.concurrent.DirectExecutorService.execute(DirectExecutorService.java:217)
at 
org.apache.flink.runtime.concurrent.FutureUtils.lambda$orTimeout$15(FutureUtils.java:582)
at 
java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at 
java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
at 
java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
at 
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at 
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

?????? flink ??????????????

2021-07-14 Thread ????????????????
??




--  --
??: 
   "user-zh"



flink??????????????

2021-11-07 Thread ??????
:kafaflink,process,
:sink??,??process??,sinkprint()
 ,???





??


 

?????? flink ??????????????

2021-11-14 Thread ????????????????
??  zk ??ha  ??D 
savepoint




--  --
??: 
   "user-zh"



Flink????????????

2021-11-14 Thread ??????????
??
    
 Flink?? 
   
 ??Flink??container??container
     

Flink????????????

2021-11-15 Thread ??????????


Flink 
WebUI??cancel??FlinkFlink??

flink????????????????????

2021-12-01 Thread ????


??????flink????????????????????

2021-12-05 Thread Xianxun Ye
flink 1.14??hybrid 
source??boundedunbounded??




??2021??12??3?? 15:21??Michael Ran ??
jdbc scan
?? 2021-12-02 14:40:06??""  ??



?????? flink????????????????????

2021-12-05 Thread ????
600w ?? 




--  --
??: 
   "user-zh"

https://ververica.github.io/flink-cdc-connectors/release-2.1/content/connectors/mysql-cdc.html


> 2021??12??2?? 2:40?? 

Flink??????????????????,????????????

2022-03-03 Thread Tony
cpu,??Flink???  32, 
512G, Flink???

?????? flink??????????????

2022-03-06 Thread ????
??




--  --
??: 
   "user-zh"



Flink????????????

2022-03-08 Thread hjw
sql??SELECT color, sum(id) FROM T GROUP BY 
colorFlinkTgroup
 by 
key??color)??Flink???

?????? Flink????????????

2022-03-08 Thread hjw
streaming api ??sql api 
streaming api




--  --
??: 
   "user-zh"

https://nightlies.apache.org/flink/flink-docs-master/docs/dev/datastream/fault-tolerance/state/

hjw <1010445...@qq.com.invalid> ??2022??3??9?? 01:32??

> sql??SELECT color, sum(id) FROM T GROUP BY
> 
colorFlinkTgroup
> by 
key??color)??Flink???

flink ????

2022-04-24 Thread ????????
flink/run   jar??

flink????

2022-07-21 Thread ??????

flink??1.14.5
??
2022-07-22 10:07:51
java.util.concurrent.TimeoutException: Heartbeat of TaskManager with id 
bdp-changlog-mid-relx-middle-promotion-dev-taskmanager-1-1 timed out.
 at 
org.apache.flink.runtime.jobmaster.JobMaster$TaskManagerHeartbeatListener.notifyHeartbeatTimeout(JobMaster.java:1299)
 at 
org.apache.flink.runtime.heartbeat.HeartbeatMonitorImpl.run(HeartbeatMonitorImpl.java:111)
 at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
 at java.util.concurrent.FutureTask.run(FutureTask.java:266)
 at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:440)
 at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:208)
 at 
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:77)
 at 
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:158)
 at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26)
 at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21)
 at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:123)
 at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)
 at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:170)
 at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
 at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
 at akka.actor.Actor$class.aroundReceive(Actor.scala:517)
 at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225)
 at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592)
 at akka.actor.ActorCell.invoke(ActorCell.scala:561)
 at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258)
 at akka.dispatch.Mailbox.run(Mailbox.scala:225)
 at akka.dispatch.Mailbox.exec(Mailbox.scala:235)
 at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
 at 
akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
 at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
 at 
akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

??

Flink

2019-02-14 Thread 龚文洲
Flink



flink????

2019-03-25 Thread IORI
ABsink,Creducesinkoperator

????: flink????

2019-03-25 Thread baiyg25...@hundsun.com
??
??AB??AC??



baiyg25...@hundsun.com
 
 IORI
?? 2019-03-26 09:46
 user-zh
?? flink
ABsink,Creducesinkoperator
 


????: ????: flink????

2019-03-25 Thread ????????qq??
DataStream ds = 


DataStream ds1 = ??ds ?? B ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 
?0?2SINK

DataStream ds2 = ??ds ?? C ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 ?0?2 
?0?2SINK


??DS1??DS2DSSQL


qq??
?0?2
?0?2baiyg25...@hundsun.com
???0?22019-03-26?0?210:09
?0?2user-zh
???0?2: flink
??
??AB??AC??
?0?2
?0?2
?0?2
baiyg25...@hundsun.com
 IORI
?? 2019-03-26 09:46
 user-zh
?? flink
ABsink,Creducesinkoperator

??????????????Flink????

2019-03-28 Thread ????

   
Flink???(Flinkjarweb??)

?????? ??????????????Flink????

2019-03-28 Thread ????





--  --
??: "Lifei Chen";
: 2019??3??29??(??) 11:10
??: "user-zh";

: Re: ??Flink



go cli, jarflink manager

https://github.com/ing-bank/flink-deployer

??

Kaibo Zhou  ??2019??3??29?? 11:08??

> ?? flink ?? Restful API ??upload  jar ?? run??
>
> ??
>
> https://ci.apache.org/projects/flink/flink-docs-stable/monitoring/rest_api.html#jars-upload
> ?? https://files.alicdn.com/tpsservice/a8d224d6a3b8b82d03aa84e370c008cc.pdf
> ??
>
>  <1010467...@qq.com> ??2019??3??28?? 9:06??
>
> > 
> >
> >
> ????Flink???(Flinkjarweb??)
>

flink????????????????

2019-04-01 Thread ????
??flink??flink??/??tm??

??????flink????????????????

2019-04-01 Thread ????
??





--  --
??: ""<1048095...@qq.com>;
: 2019??4??1??(??) 5:28
??: "user-zh";

: flink????



??flink??flink??/??tm??

??????flink????????????????

2019-04-01 Thread ????
--  --
??: ""<1048095...@qq.com>;
: 2019??4??1??(??) 5:30
??: "????";

: ??flink



??





--  --
??: ""<1048095...@qq.com>;
: 2019??4??1??(??) 5:28
??: "user-zh";

: flink



??flink??????flink??/??tm??

?????? ??????????????Flink????

2019-04-01 Thread ????

   
gitjenkinsjar??shell??JobManager??JobGraph??JobGraphJobManagerjar??
??




--  --
??: ""<1010467...@qq.com>;
: 2019??3??29??(??) 2:19
??: "user-zh";

????: ?? ??Flink








--  --
??: "Lifei Chen";
: 2019??3??29??(??) 11:10
??: "user-zh";

: Re: ??Flink????



go cli, jarflink manager

https://github.com/ing-bank/flink-deployer

??

Kaibo Zhou  ??2019??3??29?? 11:08??

> ?? flink ?? Restful API ??upload  jar ?? run??
>
> ??
>
> https://ci.apache.org/projects/flink/flink-docs-stable/monitoring/rest_api.html#jars-upload
> ?? https://files.alicdn.com/tpsservice/a8d224d6a3b8b82d03aa84e370c008cc.pdf
> ??
>
>  <1010467...@qq.com> ??2019??3??28?? 9:06??
>
> > 
> >
> >
> Flink???(Flinkjarweb??)
>

?????? ??????????????Flink????

2019-04-02 Thread ????


??lib??(??:Flink??Kafka)Flink??-C??jar??jarJobGraph??JobManager
 

StreamExecutionEnvironment.getExecutionEnvironment().getStreamGraph().getJobGraph()??JobGraphJobGraph??JobManager
   ??




--  --
??: "Yuan Yifan";
: 2019??4??2??(??) 3:24
??: "user-zh";

: Re:?? ??Flink



??env.??JAR??








?? 2019-04-02 14:39:45??"" <1010467...@qq.com> ??
>
>   
> gitjenkinsjar??shell??JobManager??JobGraph??JobGraphJobManagerjar??
>??
>
>
>
>
>--  --
>??: ""<1010467...@qq.com>;
>????: 2019??3??29??(??) 2:19
>??: "user-zh";
>
>: ?? ??Flink
>
>
>
>
>
>
>
>
>--  --
>??: "Lifei Chen";
>????: 2019??3??29??(??) 11:10
>??: "user-zh";
>
>: Re: ??Flink
>
>
>
>????go cli, jarflink manager
>
>https://github.com/ing-bank/flink-deployer
>
>??????
>
>Kaibo Zhou  ??2019??3??29?? 11:08??
>
>> ?? flink ?? Restful API ??upload  jar ?? run??
>>
>> ??
>>
>> https://ci.apache.org/projects/flink/flink-docs-stable/monitoring/rest_api.html#jars-upload
>> ?? https://files.alicdn.com/tpsservice/a8d224d6a3b8b82d03aa84e370c008cc.pdf
>> ??
>>
>>  <1010467...@qq.com> ??2019??3??28?? 9:06??
>>
>> > 
>> >
>> >
>> Flink???(Flinkjarweb??)
>>

flink??????????????

2019-04-22 Thread 1900
flink1.7.2??hadoop??2.8.5,flink on yarn 
ha?? run a job on yarn


??kafka??window??5??list??db


1.??
2.kafkatopictopic??8??8
 
3.??kafka??window
kafka??4kafka??kafka??1
??kafka
4.keyID??hash??
??kafka
??DB??kafka??
5.keyslot8??6
??2

flink??????????????

2019-05-09 Thread ??????????
 flink??

flink

2019-05-13 Thread Kobeli
hello flink watermark??
flink jobevent time?? 
source??kafka,??(under-replica), flink 
jobpartition??(kafka)??
watermark?? 
??partition??watermark
 


????: flink

2019-05-13 Thread deng
??delay ??delay flink ??

??  


?0?2
?0?2Kobeli
???0?22019-05-11?0?218:44
?0?2user-zh
???0?2flink
hello flink watermark??
flink jobevent time?? 
source??kafka,??(under-replica), flink 
jobpartition??(kafka)??
watermark?? 
??partition??watermark


flink

2019-06-13 Thread wuzongj...@hfepay.com
请问,flink写maxcomputer(odps)有相关实现案例吗?



wuzongj...@hfepay.com


flink????????

2019-06-23 Thread ????2008
??flink demo ??~~

flink????????

2019-07-21 Thread ????
??flink on yarnrabbitmq??  
rabbitmq??ideaflink1.8flink1.7


2019-07-22 11:32:12.309 [Source: Custom Source (1/1)] INFO  
org.apache.flink.runtime.taskmanager.Task  - Source: Custom Source (1/1) 
(85cfdb83f536b26e07ca2aa4a1b66302) switched from DEPLOYING
 to FAILED.
java.lang.NullPointerException: null
at 
org.apache.flink.streaming.runtime.tasks.StreamTask.createRecordWriters(StreamTask.java:1175)
at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:212)
at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:190)
at 
org.apache.flink.streaming.runtime.tasks.SourceStreamTask.(SourceStreamTask.java:51)
at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
at 
sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
at 
sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
at 
org.apache.flink.runtime.taskmanager.Task.loadAndInstantiateInvokable(Task.java:1405)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:689)
at java.lang.Thread.run(Thread.java:745)
2019-07-22 11:32:12.310 [Source: Custom Source (1/1)] INFO  
org.apache.flink.runtime.taskmanager.Task  - Freeing task resources for Source: 
Custom Source (1/1) (85cfdb83f536b26e07ca2aa4a1b663
02).
2019-07-22 11:32:12.321 [Source: Custom Source (1/1)] INFO  
org.apache.flink.runtime.taskmanager.Task  - Ensuring all FileSystem streams 
are closed for task Source: Custom Source (1/1) (85cfd
b83f536b26e07ca2aa4a1b66302) [FAILED]
2019-07-22 11:32:12.331 [flink-akka.actor.default-dispatcher-3] INFO  
org.apache.flink.runtime.taskexecutor.TaskExecutor  - Un-registering task and 
sending final execution state FAILED to Job
Manager for task Source: Custom Source 85cfdb83f536b26e07ca2aa4a1b66302.

?????? flink????????

2019-07-21 Thread ????
manager.Task  - Registering task at network: Map 
-> Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) [DEPLOYING]. 2019-07-22 05:39:02.965 
[Source: Custom Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  
- Registering task at network: Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e) [DEPLOYING]. 2019-07-22 05:39:02.987 [Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2)] INFO  
org.apache.flink.runtime.taskmanager.Task  - Map -> Filter -> 
Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2) 
(54be512ff5fc97d934ca9fbb96d66fe6) switched from DEPLOYING to RUNNING. 
2019-07-22 05:39:02.997 [Map -> Filter -> Timestamps/Watermarks -> Filter -> 
Sink: Unnamed (1/2)] INFO  org.apache.flink.runtime.taskmanager.Task  - Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) switched from DEPLOYING to RUNNING. 
2019-07-22 05:39:03.013 [Map -> Filter -> Timestamps/Watermarks -> Filter -> 
Sink: Unnamed (1/2)] INFO  org.apache.flink.streaming.runtime.tasks.StreamTask  
- No state backend has been configured, using default (Memory / JobManager) 
MemoryStateBackend (data in heap memory / checkpoints to JobManager) 
(checkpoints: 'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/checkpoints', 
savepoints: 'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/savepoints', 
asynchronous: TRUE, maxStateSize: 5242880) 2019-07-22 05:39:03.020 [Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2)] INFO  
org.apache.flink.streaming.runtime.tasks.StreamTask  - No state backend has 
been configured, using default (Memory / JobManager) MemoryStateBackend (data 
in heap memory / checkpoints to JobManager) (checkpoints: 
'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/checkpoints', savepoints: 
'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/savepoints', asynchronous: 
TRUE, maxStateSize: 5242880) 2019-07-22 05:39:03.067 [Source: Custom Source 
(1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Source: Custom Source 
(1/1) (a556fff0d64b101420a4209c8ba7be7e) switched from DEPLOYING to FAILED. 
java.lang.NullPointerException: null  at 
org.apache.flink.streaming.runtime.tasks.StreamTask.createRecordWriters(StreamTask.java:1175)
at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:212) 
 at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:190) 
 at 
org.apache.flink.streaming.runtime.tasks.SourceStreamTask.(SourceStreamTask.java:51)
   at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method) 
   at 
sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
 at 
sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
 at java.lang.reflect.Constructor.newInstance(Constructor.java:423) 
 at 
org.apache.flink.runtime.taskmanager.Task.loadAndInstantiateInvokable(Task.java:1405)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:689) 
at java.lang.Thread.run(Thread.java:745) 2019-07-22 05:39:03.068 [Source: 
Custom Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Freeing 
task resources for Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e). 2019-07-22 05:39:03.090 [Source: Custom 
Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Ensuring all 
FileSystem streams are closed for task Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e) [FAILED] 2019-07-22 05:39:03.098 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskexecutor.TaskExecutor  - Un-registering task and 
sending final execution state FAILED to JobManager for task Source: Custom 
Source a556fff0d64b101420a4209c8ba7be7e. 2019-07-22 05:39:03.176 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskexecutor.TaskExecutor  - Discarding the results 
produced by task execution a556fff0d64b101420a4209c8ba7be7e. 2019-07-22 
05:39:03.178 [flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Attempting to cancel task Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3). 2019-07-22 05:39:03.179 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Map -> Filter -> 
Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) switched from RUNNING to CANCELING. 
2019-07-22 05:39:03.179 [flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Triggering cancellation of task 
code Map -> Filter -> Timestamps/Wate

?????? flink????????

2019-07-21 Thread ????
manager.Task  - Registering task at network: Map 
-> Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) [DEPLOYING]. 2019-07-22 05:39:02.965 
[Source: Custom Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  
- Registering task at network: Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e) [DEPLOYING]. 2019-07-22 05:39:02.987 [Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2)] INFO  
org.apache.flink.runtime.taskmanager.Task  - Map -> Filter -> 
Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2) 
(54be512ff5fc97d934ca9fbb96d66fe6) switched from DEPLOYING to RUNNING. 
2019-07-22 05:39:02.997 [Map -> Filter -> Timestamps/Watermarks -> Filter -> 
Sink: Unnamed (1/2)] INFO  org.apache.flink.runtime.taskmanager.Task  - Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) switched from DEPLOYING to RUNNING. 
2019-07-22 05:39:03.013 [Map -> Filter -> Timestamps/Watermarks -> Filter -> 
Sink: Unnamed (1/2)] INFO  org.apache.flink.streaming.runtime.tasks.StreamTask  
- No state backend has been configured, using default (Memory / JobManager) 
MemoryStateBackend (data in heap memory / checkpoints to JobManager) 
(checkpoints: 'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/checkpoints', 
savepoints: 'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/savepoints', 
asynchronous: TRUE, maxStateSize: 5242880) 2019-07-22 05:39:03.020 [Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (2/2)] INFO  
org.apache.flink.streaming.runtime.tasks.StreamTask  - No state backend has 
been configured, using default (Memory / JobManager) MemoryStateBackend (data 
in heap memory / checkpoints to JobManager) (checkpoints: 
'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/checkpoints', savepoints: 
'hdfs://10.0.2.11:9000/flink/igg-flink-log-monitor/savepoints', asynchronous: 
TRUE, maxStateSize: 5242880) 2019-07-22 05:39:03.067 [Source: Custom Source 
(1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Source: Custom Source 
(1/1) (a556fff0d64b101420a4209c8ba7be7e) switched from DEPLOYING to FAILED. 
java.lang.NullPointerException: null  at 
org.apache.flink.streaming.runtime.tasks.StreamTask.createRecordWriters(StreamTask.java:1175)
at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:212) 
 at 
org.apache.flink.streaming.runtime.tasks.StreamTask.(StreamTask.java:190) 
 at 
org.apache.flink.streaming.runtime.tasks.SourceStreamTask.(SourceStreamTask.java:51)
   at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method) 
   at 
sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
 at 
sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
 at java.lang.reflect.Constructor.newInstance(Constructor.java:423) 
 at 
org.apache.flink.runtime.taskmanager.Task.loadAndInstantiateInvokable(Task.java:1405)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:689) 
at java.lang.Thread.run(Thread.java:745) 2019-07-22 05:39:03.068 [Source: 
Custom Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Freeing 
task resources for Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e). 2019-07-22 05:39:03.090 [Source: Custom 
Source (1/1)] INFO  org.apache.flink.runtime.taskmanager.Task  - Ensuring all 
FileSystem streams are closed for task Source: Custom Source (1/1) 
(a556fff0d64b101420a4209c8ba7be7e) [FAILED] 2019-07-22 05:39:03.098 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskexecutor.TaskExecutor  - Un-registering task and 
sending final execution state FAILED to JobManager for task Source: Custom 
Source a556fff0d64b101420a4209c8ba7be7e. 2019-07-22 05:39:03.176 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskexecutor.TaskExecutor  - Discarding the results 
produced by task execution a556fff0d64b101420a4209c8ba7be7e. 2019-07-22 
05:39:03.178 [flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Attempting to cancel task Map -> 
Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3). 2019-07-22 05:39:03.179 
[flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Map -> Filter -> 
Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740a6bf3006287f24c46028abcbd8ec3) switched from RUNNING to CANCELING. 
2019-07-22 05:39:03.179 [flink-akka.actor.default-dispatcher-5] INFO  
org.apache.flink.runtime.taskmanager.Task  - Triggering cancellation of task 
code Map -> Filter -> Timestamps/Watermarks -> Filter -> Sink: Unnamed (1/2) 
(740

??????flink????

2019-07-29 Thread ????
flink-conf.yaml??TaskManager??slot??




--  --
??: "";
: 2019??7??30??(??) 10:46
??: "user-zh";

: flink



??flink??job??
Could not allocate all requires slots within timeout of 30 ms. Slots 
required: 3, slots allocated: 2
flink??web??slot170??

flink????????????

2019-08-21 Thread IORI
??flink??1.7 
getExecutionEvironmentcreateRemoteEvironmentjarremoteEvironmentflink??java
 -jar
Flink??jar?? 
PackageProgram,clientclient??client,
??flink??





iPhone

????flink????????????????????????????

2019-08-26 Thread 1900
flink??flinkTwoPhaseCommitSinkFunction??,



?C beginTransaction 

?C preCommit 

?C commit 

?C abort 



sink??MYSQL
 preCommit??
commit??checkpoint??checkpoint
checkpoint??checkpoint
checkpoint??


flink????????????

2019-08-27 Thread ????
??flink??ERROR

??flink??log4j.properties??INFO
LogManager.getRootLogger().setLevel(Level.ERROR);openINFO??

?????? flink????????????

2019-08-27 Thread ????
webERROR




--  --
??: "??"; 
: 2019??8??28??(??) 1:50
??: "user-zh"; 
????: Re:?? flink




weberrores




--

??
 +86 18710107193
gaofeilong198...@163.com



?? 2019-08-27 19:51:35??"??"  ??
>??ES??error??
>
>
>
>csbl...@163.com
>Have a nice day !
>
>
>??2019??08??27?? 19:46 ??
>??ES??error??
>
>
>
>csbl...@163.com
>Have a nice day !
>
>
>??2019??08??27?? 19:29??Zili Chen ??
>?? download  ERROR ??x
>
>Best,
>tison.
>
>
> <58683...@qq.com> ??2019??8??27?? 6:56??
>
>??flink??ERROR
>
>
>??flink??log4j.properties??INFO
>
>LogManager.getRootLogger().setLevel(Level.ERROR);openINFO??

?????? ????flink????????????????????????????

2019-08-27 Thread 1900
hi,


??


public class Sink extends TwoPhaseCommitSinkFunction  {


//private Connection connection;


public Sink() {
super(new KryoSerializer <>(Connection.class , new ExecutionConfig()) , 
VoidSerializer.INSTANCE);
}


@Override
protected void invoke(Connection connection , ObjectNode objectNode , 
Context context) throws Exception {
String  stu = objectNode.get("value").toString();
Student student = JSON.parseObject(stu , Student.class);


System.err.println("start invoke..." + "id = " + student.getId() + 
"  name = " + student.getName() + "   password"
   + " = " + student.getPassword() + "  age = " + 
student.getAge());


Stringsql = "insert into Student(id,name,password,age) 
values (?,?,?,?)";
PreparedStatement ps  = connection.prepareStatement(sql);
ps.setInt(1 , student.getId());
ps.setString(2 , student.getName());
ps.setString(3 , student.getPassword());
ps.setInt(4 , student.getAge());
ps.executeUpdate();
//
if (student.getId() == 33) {
System.out.println(1 / 0);
}
}


@Override
protected Connection beginTransaction() throws Exception {
String url = "jdbc:mysql:";
return DBConnectUtil.getConnection(url , "" , "");
}


@Override
protected void preCommit(Connection connection) throws Exception {
}


@Override
protected void commit(Connection connection) {
if (connection != null) {
try {
connection.commit();
} catch (SQLException e) {
System.err.println("commit  error " + 
e.getMessage());
} finally {
try {
connection.close();
} catch (SQLException e) {
System.err.println(" finally  commit error " + 
e.getMessage());
}
}
}
}


@Override
protected void abort(Connection connection) {
if (connection != null) {
try {
connection.rollback();
} catch (SQLException e) {
System.err.println("abort error " + e.getMessage());
} finally {
try {
connection.close();
} catch (SQLException e) {
System.err.println(" finally  abort error " + 
e.getMessage());
}
}
}
    }
}







--  --
??: "Yun Tang";
: 2019??8??27??(??) 3:00
??: "user-zh";

: Re: flink



Hi

??TwoPhaseCommitSinkFunctionpreCommitsnapshotState??currentTransactionHolder??pendingCommitTransactions??notifyCheckpointComplete??commitpendingCommitTransactionspreCommit

TwoPhaseCommitSinkFunction??checkpoint 
interval??????state??



____
From: 1900 <575209...@qq.com>
Sent: Tuesday, August 27, 2019 14:15
To: user-zh 
Subject: flink

flink??flinkTwoPhaseCommitSinkFunction??,



?C beginTransaction

?C preCommit

?C commit

?C abort



sink??MYSQL
 preCommit??
commit??checkpoint??checkpoint
checkpoint??checkpoint
checkpoint??


Flink????????????

2019-09-09 Thread Evan
??flink1.7.1
??centos 7
job??start-cluster.sh
test04


??flink job??


$ bin/flink run -m test04:8081 -c 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass 
flink-scala-project1.jar
Starting execution of program


 The program finished with the following exception:

org.apache.flink.client.program.ProgramInvocationException: Could not retrieve 
the execution result. (JobID: f9ac0c76e0e44cac6d6c3b1c41afa161)
at 
org.apache.flink.client.program.rest.RestClusterClient.submitJob(RestClusterClient.java:261)
at 
org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:487)
at 
org.apache.flink.streaming.api.environment.StreamContextEnvironment.execute(StreamContextEnvironment.java:66)
at 
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1510)
at 
org.apache.flink.streaming.api.scala.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.scala:645)
at 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass$.main(FlinkConsumerKafkaSinkToKuduMainClass.scala:16)
at 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass.main(FlinkConsumerKafkaSinkToKuduMainClass.scala)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at 
sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at 
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at 
org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:529)
at 
org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)
at 
org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:427)
at 
org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:813)
at 
org.apache.flink.client.cli.CliFrontend.runProgram(CliFrontend.java:287)
at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:213)
at 
org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:1050)
at 
org.apache.flink.client.cli.CliFrontend.lambda$main$11(CliFrontend.java:1126)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at 
org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1692)
at 
org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1126)
Caused by: org.apache.flink.runtime.client.JobSubmissionException: Failed to 
submit JobGraph.
at 
org.apache.flink.client.program.rest.RestClusterClient.lambda$submitJob$8(RestClusterClient.java:380)
at 
java.util.concurrent.CompletableFuture.uniExceptionally(CompletableFuture.java:870)
at 
java.util.concurrent.CompletableFuture$UniExceptionally.tryFire(CompletableFuture.java:852)
at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474)
at 
java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977)
at 
org.apache.flink.runtime.concurrent.FutureUtils.lambda$retryOperationWithDelay$5(FutureUtils.java:216)
at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
at 
java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:736)
at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474)
at 
java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977)
at 
org.apache.flink.runtime.rest.RestClient$ClientHandler.readRawResponse(RestClient.java:515)
at 
org.apache.flink.runtime.rest.RestClient$ClientHandler.channelRead0(RestClient.java:452)
at 
org.apache.flink.shaded.netty4.io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:340)
at 
org.apache.flink.shaded.netty4.io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:286)
at 
org.apache.flink.shaded.netty4

Flink????????????

2019-09-09 Thread Evan
??flink1.7.1
??centos 7
job??start-cluster.sh
test04


??flink job??


$ bin/flink run -m test04:8081 -c 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass 
flink-scala-project1.jar
Starting execution of program


 The program finished with the following exception:

org.apache.flink.client.program.ProgramInvocationException: Could not retrieve 
the execution result. (JobID: f9ac0c76e0e44cac6d6c3b1c41afa161)
at 
org.apache.flink.client.program.rest.RestClusterClient.submitJob(RestClusterClient.java:261)
at 
org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:487)
at 
org.apache.flink.streaming.api.environment.StreamContextEnvironment.execute(StreamContextEnvironment.java:66)
at 
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1510)
at 
org.apache.flink.streaming.api.scala.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.scala:645)
at 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass$.main(FlinkConsumerKafkaSinkToKuduMainClass.scala:16)
at 
org.apache.flink.quickstart.factory.FlinkConsumerKafkaSinkToKuduMainClass.main(FlinkConsumerKafkaSinkToKuduMainClass.scala)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at 
sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at 
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at 
org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:529)
at 
org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)
at 
org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:427)
at 
org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:813)
at 
org.apache.flink.client.cli.CliFrontend.runProgram(CliFrontend.java:287)
at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:213)
at 
org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:1050)
at 
org.apache.flink.client.cli.CliFrontend.lambda$main$11(CliFrontend.java:1126)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at 
org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1692)
at 
org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1126)
Caused by: org.apache.flink.runtime.client.JobSubmissionException: Failed to 
submit JobGraph.
at 
org.apache.flink.client.program.rest.RestClusterClient.lambda$submitJob$8(RestClusterClient.java:380)
at 
java.util.concurrent.CompletableFuture.uniExceptionally(CompletableFuture.java:870)
at 
java.util.concurrent.CompletableFuture$UniExceptionally.tryFire(CompletableFuture.java:852)
at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474)
at 
java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977)
at 
org.apache.flink.runtime.concurrent.FutureUtils.lambda$retryOperationWithDelay$5(FutureUtils.java:216)
at 
java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
at 
java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:736)
at 
java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474)
at 
java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977)
at 
org.apache.flink.runtime.rest.RestClient$ClientHandler.readRawResponse(RestClient.java:515)
at 
org.apache.flink.runtime.rest.RestClient$ClientHandler.channelRead0(RestClient.java:452)
at 
org.apache.flink.shaded.netty4.io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
at 
org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:340)
at 
org.apache.flink.shaded.netty4.io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:286)
at 
org.apache.flink.shaded.netty4

flink????????????????

2019-12-17 Thread cs
flink run -m yarn-cluster -yn 10 -ys 1 -p 10 -yn 20 -ys 2 -p 
40 


????flink??????????????????????

2019-12-24 Thread 1530130567

  ??flink stream api??ETL
  
1??tumble??watermark??10s??
  topic 5000/s,??topic4000/s
  
processfunction
  ps

?????? ????flink??????????????????????

2019-12-24 Thread 1530130567

??metricrecordsIn>recordsOut
??window??processfunction??filter??
??
.window(TumblingEventTimeWindows.of(org.apache.flink.streaming.api.windowing.time.Time.minutes(1)))
.process(new ProcessWindowFunction

flink??????????????

2019-12-31 Thread cs
taskmanager15Gheap10G
tm??1.cutoff(15GB * 0.25) 
2.heap(heap15GB - cutoff) 
3.offheap(offheap??15GB-heap)
offheap??-XX:MaxDirectMemorySize??
MaxDirectMemorySize??

?????? flink??????????????

2019-12-31 Thread cs
??
??tm-XX:NewSize
tm15G -XX:NewSize=2442764288
tm20G ?? -XX:NewSize=2442764288
??




--  --
??: "Xintong Song"

????????flink??????????????

2020-01-16 Thread sun
flink??redisredis??flinkflink??flink??redis

flink????????????????

2020-03-19 Thread 512348363
??
DataStream

??????????flink????????????

2020-03-23 Thread ????
??es??flinkhdfssetBatchSiz
 ??hive

flink ???????? ????????

2020-04-05 Thread ??????
2020-04-04 18:58:41.515 ERROR 1 --- [ent-IO-thread-4] 
org.apache.flink.runtime.rest.RestClient.parseResponse:393 : Received response 
was neither of the expected type ([simple type, class 
org.apache.flink.runtime.rest.messages.job.JobExecutionResultResponseBody]) nor 
an error. 
Response=JsonResponse{json={"status":{"id":"COMPLETED"},"job-execution-result":{"id":"95bfafbc9b8e31565d44ae0018a0af7d","application-status":"FAILED","accumulator-results":{},"net-runtime":3098,"failure-cause":{"class":"java.lang.OutOfMemoryError","stack-trace":"java.lang.OutOfMemoryError:
 Java heap 
space\n","serialized-throwable":"rO0ABXNyAClvcmcuYXBhY2hlLmZsaW5rLnV0aWwuU2VyaWFsaXplZFRocm93YWJsZWUWnfUfpxPzAgADTAAZZnVsbFN0cmluZ2lmaWVkU3RhY2tUcmFjZXQAEkxqYXZhL2xhbmcvU3RyaW5nO0wAFm9yaWdpbmFsRXJyb3JDbGFzc05hbWVxAH4AAVsAE3NlcmlhbGl6ZWRFeGNlcHRpb250AAJbQnhyABNqYXZhLmxhbmcuRXhjZXB0aW9u0P0fPho7HMQCAAB4cgATamF2YS5sYW5nLlRocm93YWJsZdXGNSc5d7jLAwAETAAFY2F1c2V0ABVMamF2YS9sYW5nL1Rocm93YWJsZTtMAA1kZXRhaWxNZXNzYWdlcQB+AAFbAApzdGFja1RyYWNldAAeW0xqYXZhL2xhbmcvU3RhY2tUcmFjZUVsZW1lbnQ7TAAUc3VwcHJlc3NlZEV4Y2VwdGlvbnN0ABBMamF2YS91dGlsL0xpc3Q7eHBwdAAPSmF2YSBoZWFwIHNwYWNldXIAHltMamF2YS5sYW5nLlN0YWNrVHJhY2VFbGVtZW50OwJGKjw8/SI5AgAAeHAAc3IAJmphdmEudXRpbC5Db2xsZWN0aW9ucyRVbm1vZGlmaWFibGVMaXN0/A8lMbXsjhACAAFMAARsaXN0cQB+AAd4cgAsamF2YS51dGlsLkNvbGxlY3Rpb25zJFVubW9kaWZpYWJsZUNvbGxlY3Rpb24ZQgCAy173HgIAAUwAAWN0ABZMamF2YS91dGlsL0NvbGxlY3Rpb247eHBzcgATamF2YS51dGlsLkFycmF5TGlzdHiB0h2Zx2GdAwABSQAEc2l6ZXhwAHcEAHhxAH4AEXh0ACxqYXZhLmxhbmcuT3V0T2ZNZW1vcnlFcnJvcjogSmF2YSBoZWFwIHNwYWNlCnQAGmphdmEubGFuZy5PdXRPZk1lbW9yeUVycm9ydXIAAltCrPMX+AYIVOACAAB4cf6s7QAFc3IAGmphdmEubGFuZy5PdXRPZk1lbW9yeUVycm9ycjG7cIiI4xUCAAB4cgAdamF2YS5sYW5nLlZpcnR1YWxNYWNoaW5lRXJyb3I5wlZUgC8OHgIAAHhyAA9qYXZhLmxhbmcuRXJyb3JFHTZWi4IOVgIAAHhyABNqYXZhLmxhbmcuVGhyb3dhYmxl1cY1Jzl3uMsDAARMAAVjYXVzZXQAFUxqYXZhL2xhbmcvVGhyb3dhYmxlO0wADWRldGFpbE1lc3NhZ2V0ABJMamF2YS9sYW5nL1N0cmluZztbAApzdGFja1RyYWNldAAeW0xqYXZhL2xhbmcvU3RhY2tUcmFjZUVsZW1lbnQ7TAAUc3VwcHJlc3NlZEV4Y2VwdGlvbnN0ABBMamF2YS91dGlsL0xpc3Q7eHBwdAAPSmF2YSBoZWFwIHNwYWNldXIAHltMamF2YS5sYW5nLlN0YWNrVHJhY2VFbGVtZW50OwJGKjw8/SI5AgAAeHABc3IAG2phdmEubGFuZy5TdGFja1RyYWNlRWxlbWVudGEJxZomNt2FAgAESQAKbGluZU51bWJlckwADmRlY2xhcmluZ0NsYXNzcQB+AAVMAAhmaWxlTmFtZXEAfgAFTAAKbWV0aG9kTmFtZXEAfgAFeHCAdAAAcHEAfgAOcHg="}}},
 httpResponseStatus=200 OK}


org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException:
 Unrecognized field "status" (class 
org.apache.flink.runtime.rest.messages.ErrorResponseBody), not marked as 
ignorable (one known property: "errors"])
 at [Source: N/A; line: -1, column: -1] (through reference chain: 
org.apache.flink.runtime.rest.messages.ErrorResponseBody["status"])
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException.from(UnrecognizedPropertyException.java:62)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.DeserializationContext.reportUnknownProperty(DeserializationContext.java:851)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.std.StdDeserializer.handleUnknownProperty(StdDeserializer.java:1085)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.handleUnknownProperty(BeanDeserializerBase.java:1392)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.handleUnknownProperties(BeanDeserializerBase.java:1346)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializer._deserializeUsingPropertyBased(BeanDeserializer.java:455)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.deserializeFromObjectUsingNonDefault(BeanDeserializerBase.java:1127)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserializeFromObject(BeanDeserializer.java:298)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserialize(BeanDeserializer.java:133)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper._readValue(ObjectMapper.java:3779)
 ~[flink-shaded-jackson-2.7.9-6.0.jar!/:2.7.9-6.0]
at 
org.ap

?????? flink ???????? ????????

2020-04-06 Thread ??????
??
  
??jvm??
 flink on k8s, flink 
job,??k8s??




--  --
??: "Xintong Song"https://ci.apache.org/projects/flink/flink-docs-release-1.10/zh/ops/memory/mem_trouble.html#outofmemoryerror-java-heap-space

Thank you~

Xintong Song



On Mon, Apr 6, 2020 at 12:33 AM ?? 

flink ????????????????????????????????????

2020-04-26 Thread ????
HI ALL ??
     flink
           /abc/202004*/t1.data  
??2020??4??t1.data??
           /abc/20200401/t*.data 
??2020??4??1t??
     ??

?????? flink ????????????????????????????????????

2020-04-27 Thread ????
??
DatasetFileInputFormat??supportsMultiPaths??Deprecated
/**
 * Override this method to supports multiple paths.
 * When this method will be removed, all FileInputFormats have to support 
multiple paths.
 *
 * @return True if the FileInputFormat supports multiple paths, false otherwise.
 *
 * @deprecated Will be removed for Flink 2.0.
 */
@Deprecated
public boolean supportsMultiPaths() {
   return false;
}




--  --
??: "Jingsong Li"

flink ????????????

2020-06-08 Thread ??????
??
flink??
1.table1.insertsink_table
2.sink_table.insertsink_table1
??flinkflink??

?????? flink ????????????

2020-06-08 Thread 1048262223
Hi


Flink 
source??sink


source(source1) -> transform -> sink(sink1)


source(sink1) -> transform -> sink(sink2)





Best,
Yichao Yang




--  --
??: "??"<201782...@qq.com>;
: 2020??6??8??(??) 5:52
??: "user-zh"

Flink??????????????????

2020-06-08 Thread Z-Z
Hi?? ??
    
??Flink??(NullPointer??)checkpoint??savepoint??
    1:  Flink??
    
2??checkpoint??savepointsavepoint??
    
3??

??????Flink??????????????????

2020-06-08 Thread 1048262223
Hi


1.try catch??
2.??
3.??try catch + ??perf log + 


Best,
Yichao Yang




--  --
??: "Z-Z"

??????flink????????????????

2020-06-09 Thread 1048262223
Hi


rich functionopenbroadcast??


Best,
Yichao Yang




--  --
??: "zjfpla...@hotmail.com"

????: ??????flink????????????????

2020-06-09 Thread zjfpla...@hotmail.com
??



zjfpla...@hotmail.com
 
 1048262223
?? 2020-06-09 20:17
 user-zh
?? ??flink
Hi
 
 
rich functionopenbroadcast??
 
 
Best,
Yichao Yang
 
 
 
 
--  --
??: "zjfpla...@hotmail.com"

?????? ??????flink????????????????

2020-06-09 Thread 1048262223
Hi


Broadcast[1]??open??open??




[1]https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/state/broadcast_state.html


Best,
Yichao Yang


--  --
??: "zjfpla...@hotmail.com"

flink??????????????????

2020-06-09 Thread ??????
>Hi,
>??flink??kafka??EXACTLY-ONCE??
>??debug??invoketraction.producer.send()??precommit??commit
>??

??????flink??????????????????

2020-06-09 Thread Yichao Yang
Hi


sink  
??kafkakafka1.0??kafkaEXACTLY-ONCE


Best,
Yichao Yang




--  --
??: "??"

??????flink??????????????????

2020-06-10 Thread 1193216154
read uncommit ??read commint.
read uncommit ??flink??commit??
read commit??




--  --
??: "??"

??????flink????????????????

2020-06-11 Thread Yichao Yang
Hi


broadcast??


Best,
Yichao Yang




--  --
??: "xue...@outlook.com"https://go.microsoft.com/fwlink/?LinkId=550986>;

??????flink??????????????????

2020-06-11 Thread ??????
>Hi
>
>exctly-once
>checkpoint
>kafka
>??






--  --
??: "Matt Wang"

  1   2   3   4   5   6   7   8   9   10   >