Re: Flink connector 是否支持忽略delete message

2023-07-10 Thread yh z
Hi, shi franke. 你可以尝试自己实现一个 DynamicTableSink,在里面添加参数 “sink.ignore-delete”。 你可以参考 github 上的一些实现,例如 clickhouse: https://github.com/liekkassmile/flink-connector-clickhouse-1.13 shi franke 于2023年7月7日周五 19:24写道: > >

Re: flink1.17.1使用kafka source异常

2023-07-05 Thread yh z
Hi, aiden. 看起来是类冲突,按照官方的文档,使用 kafka 时,你应该是不需要引入 flink-core 和 flink-connector-base 的( https://nightlies.apache.org/flink/flink-docs-release-1.17/docs/connectors/table/kafka/)。如果是因为其他原因需要使用这两个 jar, 你可以使用 mvn dependency::tree 查看一下 "org/apache/kafka/clients/consumer/ConsumerRecord"

Re: Flink 1.16 流表 join 的 FilterPushDown 及并行

2023-07-05 Thread yh z
Hi, Chai Kelun, 你的 filter condition 里面包含了你自定义的 UDF,是不满足 filter push down 的条件的,因为对于优化器来说 UDF 是不确定的,优化器不能从里面提取到可以下推的条件, 如果你想实现下推,可以尝试抽取下确定性的 condition,如 product.id > 10 etc.。另外,Flink 是支持 broadcast hash join 的,如果你想控制某两个表的 join type,你可以通过 join hint 来指定 join 类型为 broadcast。() Chai Kelun 于2023年7月3日周一

Setting up a Multi-node Flink Cluster

2023-06-06 Thread Z M Ang
Hello Is there a reference working implementation of a Multi-VM Flink Cluster (NOT on Docker)? e.g., a 1 Master VM + 3 Worker VMs. Not looking for documentation - but a working example with conf files modified. Thanks Z Mang

Flink job submission to Multi-VM Flink Cluster fails!

2023-06-02 Thread Z M Ang
Hello, I can launch a Flink cluster (version 1.17.x) on my laptop with 1 Job Manager and 3 Task Managers. The cluster starts, jobs can be submitted correctly on the localhost (my laptop). Next I tried to launch this cluster on 4 VMs - 1 Master

????

2023-03-31 Thread z

Re: flink cdc能否同步DDL语句?

2022-10-10 Thread yh z
目前,社区的 Flink CDC,是可以读取到 binlog 中的 DDL 语句的,但是不会传递给下游,即下游无法感知 schema 的变化。 Xuyang 于2022年10月10日周一 16:46写道: > Hi, 目前应该是不行的 > 在 2022-09-26 23:27:05,"casel.chen" 写道: > >flink cdc能否同步DDL语句以实现schema同步? 例如create table , create index, truncate > table等 >

Re: flink 的数据传输,是上游算子推给下游, 还是下游算子拉取上游, 设计的考虑是啥?

2022-09-20 Thread yh z
你好。 Flink 采用的是 pull 模型。pull 模型的优点在于:1. 其具有更好的扩展性(下游的消费者可以根据需求增加,只需要获取到上游的消费位点); 2. 下游的消费者可以根据需求来调整消费速率; 3.网络传输,flink 以前也尝试使用过push模型,且为了节约开销,进程间是复用 TCP连接,一个 task 线程的性能瓶颈将导致整条链路的所有 task 线程不能接收数据,影响整体的数据消费速率。 push模型的优点:消耗较小,不需要设计机制来一直轮训观察上游节点的数据情况。 Xuyang 于2022年9月9日周五 20:35写道: >

Re: flink hybrid source问题

2022-09-20 Thread yh z
你好,hybrid source 现在需要基于 FLIP-27 source 来实现(如:FileSource, KafkaSource),对于非 FLIP-27 source 需要做一些修改后才可以使用。如果想参与 hybird source 的扩展,可以在 slack 中加入flink社群,并发起讨论。 关于 source 相关的文档,可以查看官网和 flip 设计和讨论页面( https://cwiki.apache.org/confluence/display/FLINK/FLIP-150%3A+Introduce+Hybrid+Source )(

Re: 这里为什么会报null指针错误,和源表数据有关系吗?

2022-09-18 Thread yh z
Hi,看起来是需要一个脏数据处理的逻辑的,可以在jira上建一个issue,并给出更加详细的报错信息和可能的脏数据信息呀。建立issue后,可以在邮件里回复下。 Asahi Lee 于2022年9月14日周三 09:33写道: > 是的,运行一段时间后发生错误,null指针的代码哪里是不是要非空处理? > > > > > --原始邮件-- > 发件人: > "user-zh" >

Re: 关于flink table store的疑问

2022-09-13 Thread yh z
你好,从我个人的角度出发,我认为 flink-table-store 与 hudi, iceberg 的定位是不同的。 hudi 和 iceberg 更多的是一种 format 格式,通过这个格式来管理 schema 信息和解决行业痛点,其不与特定计算引擎绑定。其中, hudi 解决了超大数据量下的 upsert 问题, iceberg 解决了 oss 存储和上云的问题,但是他们本质上还是一种存储格式(format),这是其优势也是其劣势,优势在于不受引擎约束,专注于format层本身;缺点是无法参与主流引擎的未来规划,不易扩展,且发展受限,不能很快的参与到 olap等领域。 而

Re: 触发savepoint后, source算子会从对应offset处停止消费吗?

2022-09-08 Thread yh z
hi, 在我的理解里,savePoint 的作用和 checkPoint 是类似的,只是在 flink 1.16 以前 savePoint 只支持全量的 savePoint,底层都是采用的 barrier 实现机制。但是在1.16的规划文档里( https://cwiki.apache.org/confluence/display/FLINK/FLIP-203%3A+Incremental+savepoints),savepoint 也将支持增量的模式。 当 savepoint 触发时, source 会去保存状态,是会停止消费的。 郑 致远 于2022年9月8日周四 19:39写道:

Re: hello flink

2022-09-02 Thread yh z
Hello yh z 于2022年9月2日周五 11:51写道: > hello flink >

s3p 如果在本地调试

2022-05-19 Thread z y xing
各位好: 了解实际运行是要复制jar到plugin下,但是调试的话用该怎么初始化s3p这个文件系统了? flink版本 1.14,win10 项目通过flink-quick-start创建,在pom中添加了如下依赖 org.apache.flink flink-s3-fs-presto ${flink.version} 初始代码类似如下: Configuration fileSystemConf = new Configuration();

flink sql????????????

2021-09-28 Thread z
hi??kafkaflink sqlmysql??Aid??tsjoin??

动态加载table和udf的方法

2020-10-09 Thread Zeahoo Z
你好,在开发中遇到了下面这个困难。 目前将定义的table和udf函数写在了 conf/sql-client-defaults.yaml 文件中。程序运行没有问题,但是带来一个问题:如果我需要添加或者修改table或者udf文件的时候,需要重启flink程序。有没有办法可以让我动态地添加或者更新我的table和udf(table新增可以通过sql-client来添加,但是更倾向于将table的定义记录在文件,udf存在修改和新增的情况)。这样依赖可以保证flink不重启。

??????flink 1.11 taskmanager????????????????????????

2020-09-09 Thread Z-Z
rocksdb?? ---- ??: "Z-Z"

flink 1.11 taskmanager????????????????????????

2020-09-09 Thread Z-Z
Hi ?? flink docker sessiontaskmanager taskmanager.memory.process.size: 5120m taskmanager.memory.jvm-metaspace.size: 1024m taskmanager7.5G??taskmanager INFO [] - Final TaskExecutor Memory

?????? Flink Cli ????????

2020-07-20 Thread Z-Z
cli ?? ??jobmanager bin/flink run -d -p 1 -s {savepointuri} /data/test.jar webui??http://jobmanager:8081; submit new job??jar??savepoint path

?????? Flink Cli ????????

2020-07-19 Thread Z-Z
taskmanagercliwebui?? 2020-07-20 03:29:25,959 WARN org.apache.kafka.clients.consumer.ConsumerConfig - The configuration 'value.serializer' was supplied but isn't a known config. 2020-07-20 03:29:25,959 INFO

?????? Flink Cli ????????

2020-07-17 Thread Z-Z
Flink 1.10.0 ,taskmanager?? 2020-07-17 15:06:43,913 ERROR org.apache.flink.contrib.streaming.state.RocksDBKeyedStateBackendBuilder - Caught unexpected exception. java.io.EOFException at java.io.DataInputStream.readFully(DataInputStream.java:197) at

Flink Cli ????????

2020-07-17 Thread Z-Z
restAPI??savepoint(??/jobs/overview --- /jobs/{jobid}/savepoints --- /jobs/{jobid}/savepoints/{triggerid})??flinksavepointwebuijar??savepoint??

????????Flink Hadoop????????

2020-07-14 Thread Z-Z
Flink 1.11.0docker-compose??docker-compose?? version: "2.1" services: jobmanager: image: flink:1.11.0-scala_2.12 expose: - "6123" ports: - "8081:8081" command: jobmanager environment: - JOB_MANAGER_RPC_ADDRESS=jobmanager -

Flink Hadoop????

2020-07-07 Thread Z-Z
Hi?? ??Flink 1.10.0??jobmanager??libflink-shaded-hadoop-2-uber-2.7.5-10.0.jar??webuicli??Could not find a file system implementation for scheme 'hdfs'. The scheme is not directly supported by Flink

Flink????????

2020-07-06 Thread Z-Z
Hi?? ?? Flink??checkpointcheckpoint??

keyed state????????????????

2020-06-10 Thread Z-Z
??keybyStateDescriptorkeyed state??

???? StreamingFileSink????????????????HA??Hadoop????,??????yarn job

2020-06-09 Thread ???Z?w???w
StreamingFileSinkHA??Hadoop,??yarn job ?? ?? ?? = Mobile??18611696624 QQ:79434564

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

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

Re: Flink cannot recognized catalog set by registerCatalog.

2019-08-12 Thread Xuefu Z
nStreamingMode() > .withBuiltInCatalogName("ca1") > .withBuiltInDatabaseName("db1") > .build()); > > > As Dawid said, if I want to store in my custom catalog, I can call > catalog.createTable or using DDL. > > Thanks, > SImon > > On

Re: Flink cannot recognized catalog set by registerCatalog.

2019-08-12 Thread Xuefu Z
Hi Simon, Thanks for reporting the problem. There is some rough edges around catalog API and table environments, and we are improving post 1.9 release. Nevertheless, tableEnv.registerCatalog() is just to put a new catalog in Flink's CatalogManager, It doens't change the default catalog/database

Re: [ANNOUNCE] Rong Rong becomes a Flink committer

2019-07-11 Thread Xuefu Z
Congratulations, Rong! On Thu, Jul 11, 2019 at 10:59 AM Bowen Li wrote: > Congrats, Rong! > > > On Thu, Jul 11, 2019 at 10:48 AM Oytun Tez wrote: > > > Congratulations Rong! > > > > --- > > Oytun Tez > > > > *M O T A W O R D* > > The World's Fastest Human Translation Platform. > >

re: How can I submit a flink job to YARN/Cluster from java code?

2017-07-30 Thread z...@zjdex.com
You can see the detail information at https://ci.apache.org/projects/flink/flink-docs-release-1.3/monitoring/rest_api.html. In this link, it is not said the restful API can not used in yarn cluster. I think it is ok, and you can try. :) z...@zjdex.com ?? 2017-07-28

re: How can I submit a flink job to YARN/Cluster from java code?

2017-07-28 Thread z...@zjdex.com
t_api.html . z...@zjdex.com ?? 2017-07-28 15:50 user ?? How can I submit a flink job to YARN/Cluster from java code? Hello?? I want to submit a flink job to YARN/Cluster from java code.If this is feasible? Is there anyone tried to do it before ?? Thanks