Apache Flink是一个用于可扩展批处理和流数据处理的开源平台。 Flink在一个系统中支持批量和流分析。分析程序可以用Java和Scala中简洁优雅的API编写。
在不使用 Thread.sleep 的情况下限制 Flink 作业中的事件流动
我是 Flink 新手,我正在尝试实现一个从 Kafka 主题消费的管道,对该数据执行较小的过滤和转换,并异步写入端点。
我目前正在编写我的第一个 Flink 应用程序,并且想要监视文件夹中的新文件。不幸的是我找不到关于这个主题的很多例子。 我找到了 readFile(fileInputFormat,
我们正在尝试按我们从三个 Kafka 主题消费的事件时间对事件进行排序。每个源主题都有三个分区,我们将 Flink 并行度也设置为 3。阅读事件后...
无法在我的 Maven 构建中获取 org.apache.flink.formats 包
我一直在尝试通过maven构建我的apache flink项目,但由于某种原因我遇到了编译错误。值得注意的是“org.apache.flink.formats 包不存在”...
我有一个flink应用程序,我从两个kafka源读取数据并对两个流执行连接操作 我在源头设置水印策略,例如 数据流...
无法将 Flink SQL 作业升级到 1.18,因为 Calc 和 ChangelogNormalize 顺序发生了变化
上下文 我们在使用版本 1.15.2 的 Flink 集群中运行 Flink 作业。这些工作包括: 一个或多个 KafkaSource,我们从中创建带有主键的变更日志流。 SQL
java.lang.reflect.InaccessibleObjectException:模块java.base不会向未命名模块“打开java.util.concurrent.atomic”
我正在编写一个 apache flink 程序来本地运行并与 google pubsub 交互。 依赖关系 17 ...
我正在尝试在 Flink 中的 KeyedStream 上执行映射操作: Stream.map(new JsonToMessageObjectMapper()) .keyBy("关键字段") .map(新的消息处理器St...
如何通过保存点恢复作业,并在应用程序模式下通过executeAsync运行多作业(flink 1.18)
我正在开发 1.18 版本的 flink java,并希望使用应用程序模式在一个 pod 中运行 2 个作业(k8s docker 部署)。 在java代码中,我使用for语句使用env创建2个或更多作业。
即使所有分区中都没有新数据,如何在 Flink SQL 中为 Kafka 源提前水印?
我有一个简单的Flink(v1.17)SQL流作业,使用Kafka作为源,我配置了一些与水印相关的配置,但我似乎不明白如何强制
在 Apache Flink DataStream API 中访问 Kafka 元数据
使用 Apache Flink DataStream Connectors for Kafka 创建源时如何访问 Kafka 元数据?,我注意到在 Apache Flink Table API Connector for Kafka 中,我们能够访问记录
我想使用 Flink 合并两个(多个)流。两个流本身都是有序的,我希望合并结果也被排序。举个例子 [1,2,4,5,7,8,...] 和 [2,3,6,7,..] 应该产生
例如,我有一个 flink 窗口应用程序,其数据源是一个发送来自不同公司的员工姓名的 kafka 流。现在,如果一家公司的数据源停止提取数据,我想要
Flink 似乎有数据采样器,但我没有看到任何如何在 Flink 流应用程序中应用数据采样器的示例。 是否可以在 Flink Streaming 中应用数据采样器
我有一个 Flink 作业,我不明白为什么它不会打印到标准输出。我注意到,如果删除过滤器和水印,我会看到来自我的 kafka 主题的原始消息。但是应用聚合...
Flink Timer onTimer 事件 - 它们是否强制重新分配流?
我有 KeyedProcessFunction (名为 Stateless),可以处理每个键的数据(例如事务 id)。 无状态是轻量级的,不使用状态——它只设置计时器一个计时器。 我正在上游阅读
无法使用 apache flink 的接收器功能将带有标头的 CSV 文件写入 S3 存储桶
我的项目需要使用 apache flink(版本 1.18.0)的接收器功能将 csv 文件写入 S3 存储桶,其中包含标头。使用的编程语言是java。 Hadoop 文件系统是...
为什么Flink Exactly Once commit不会失败?
我正在使用 Flink EO。重新启动 Flink 作业时,我收到此警告: 2023-12-30 13:07:44.538 [Co-Flat 地图 -> Sink: Sink1 (3/8)#0] 警告 o.a.f.s.api.functions.sink.TwoPhaseCommitSinkFunc...
为什么 Flink Exactly Once commit 不会失败?
我正在使用 Flink EO。重新启动 Flink 作业时,我收到此警告: 2023-12-30 13:07:44.538 [Co-Flat 地图 -> Sink: Sink1 (3/8)#0] 警告 o.a.f.s.api.functions.sink.TwoPhaseCommitSinkFunc...
我在使用 Apache Flink 写入 HBase 表时遇到问题。我已经成功配置了从 Kafka 读取和写入以及从 HBase 读取 RowKey。然而,当尝试...