首页/新闻资讯/正文详情

【flink番外篇】18、通过数据管道将table source加入datastream示例

发布时间:2026/9/26 2:00:25 来源:云帆数科 栏目:资讯中心
【flink番外篇】18、通过数据管道将table source加入datastream示例
最近在研究 AI BI智能数据分析 的落地实践。敬请期待后续专题实战系列《从零手把手教你搭建 AI 驱动的 BI 系统》将覆盖 Text2SQL、多轮对话、语义层、权限治理、生产级部署全链路代码可落地、坑点全复盘。一、Flink 专栏Flink 专栏系统介绍某一知识点并辅以具体的示例进行说明。1、Flink 部署系列本部分介绍Flink的部署、配置相关基础内容。2、Flink基础系列本部分介绍Flink 的基础部分比如术语、架构、编程模型、编程指南、基本的datastream api用法、四大基石等内容。3、Flik Table API和SQL基础系列本部分介绍Flink Table Api和SQL的基本用法比如Table API和SQL创建库、表用法、查询、窗口函数、catalog等等内容。4、Flik Table API和SQL提高与应用系列本部分是table api 和sql的应用部分和实际的生产应用联系更为密切以及有一定开发难度的内容。5、Flink 监控系列本部分和实际的运维、监控工作相关。二、Flink 示例专栏Flink 示例专栏是 Flink 专栏的辅助说明一般不会介绍知识点的信息更多的是提供一个一个可以具体使用的示例。本专栏不再分目录通过链接即可看出介绍的内容。两专栏的所有文章入口点击Flink 系列文章汇总索引文章目录一、DataStream 和 Table集成-数据管道1、maven依赖2、Adding Table API Pipelines to DataStream API 示例本文介绍了将table api管道加入datastream。如果需要了解更多内容可以在本人Flink 专栏中了解更新系统的内容。本文除了maven依赖外没有其他依赖。更多详细内容参考文章21、Flink 的table API与DataStream API 集成完整版一、DataStream 和 Table集成-数据管道1、maven依赖propertiesencodingUTF-8/encodingproject.build.sourceEncodingUTF-8/project.build.sourceEncodingmaven.compiler.source1.8/maven.compiler.sourcemaven.compiler.target1.8/maven.compiler.targetjava.version1.8/java.versionscala.version2.12/scala.versionflink.version1.17.0/flink.version/propertiesdependenciesdependencygroupIdorg.apache.flink/groupIdartifactIdflink-clients/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-java/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-common/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-streaming-java/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-api-java-bridge/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-sql-gateway/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-csv/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-json/artifactIdversion${flink.version}/versionscopeprovided/scope/dependency!-- https://mvnrepository.com/artifact/org.apache.flink/flink-table-planner --dependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-planner_2.12/artifactIdversion${flink.version}/versionscopeprovided/scope/dependency!-- https://mvnrepository.com/artifact/org.apache.flink/flink-table-api-java-uber --dependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-api-java-uber/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-runtime/artifactIdversion${flink.version}/versionscopeprovided/scope/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-connector-jdbc/artifactIdversion3.1.0-1.17/version/dependencydependencygroupIdmysql/groupIdartifactIdmysql-connector-java/artifactIdversion5.1.38/version/dependency!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-hive --dependencygroupIdorg.apache.flink/groupIdartifactIdflink-connector-hive_2.12/artifactIdversion1.17.0/version/dependencydependencygroupIdorg.apache.hive/groupIdartifactIdhive-exec/artifactIdversion3.1.2/version/dependency!-- flink连接器 --!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-kafka --dependencygroupIdorg.apache.flink/groupIdartifactIdflink-connector-kafka/artifactIdversion${flink.version}/version/dependency!-- https://mvnrepository.com/artifact/org.apache.flink/flink-sql-connector-kafka --dependencygroupIdorg.apache.flink/groupIdartifactIdflink-sql-connector-kafka/artifactIdversion${flink.version}/versionscopeprovided/scope/dependency!-- https://mvnrepository.com/artifact/org.apache.commons/commons-compress --dependencygroupIdorg.apache.commons/groupIdartifactIdcommons-compress/artifactIdversion1.24.0/version/dependencydependencygroupIdorg.projectlombok/groupIdartifactIdlombok/artifactIdversion1.18.2/version!-- scopeprovided/scope --/dependency/dependencies2、Adding Table API Pipelines to DataStream API 示例单个Flink作业可以由多个相邻运行的断开连接的管道组成。Table API中定义的Source-to-sink管道可以作为一个整体附加到StreamExecutionEnvironment并在调用DataStream API中的某个执行方法时提交。源不一定是table source也可以是以前转换为Table API的另一个DataStream管道。因此可以将 table sinks用于DataStream API程序。通过使用StreamTableEnvironment.createStatementSet()创建的专用StreamStatementSet实例可以使用该功能。通过使用语句集planner 可以一起优化所有添加的语句并在调用StreamStatement set.attachAsDataStream()时提供一个或多个添加到StreamExecutionEnvironment的端到端管道( end-to-end pipelines)。下面的示例演示如何将表程序添加到一个作业中的DataStream API程序。importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.api.functions.sink.DiscardingSink;importorg.apache.flink.table.api.DataTypes;importorg.apache.flink.table.api.Schema;importorg.apache.flink.table.api.Table;importorg.apache.flink.table.api.TableDescriptor;importorg.apache.flink.table.api.bridge.java.StreamStatementSet;importorg.apache.flink.table.api.bridge.java.StreamTableEnvironment;/** * author alanchan * */publicclassTestTablePipelinesToDataStreamDemo{/** * param args * throws Exception */publicstaticvoidmain(String[]args)throwsException{StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();StreamTableEnvironmenttenvStreamTableEnvironment.create(env);StreamStatementSetstatementSettenv.createStatementSet();// 建立数据源TableDescriptorsourceDescriptorTableDescriptor.forConnector(datagen).option(number-of-rows,3).schema(Schema.newBuilder().column(myCol,DataTypes.INT()).column(myOtherCol,DataTypes.BOOLEAN()).build()).build();// 建立sinkTableDescriptorsinkDescriptorTableDescriptor.forConnector(print).build();// add a pure Table API pipelineTabletableFromSourcetenv.from(sourceDescriptor);statementSet.add(tableFromSource.insertInto(sinkDescriptor));// use table sinks for the DataStream API pipelineDataStreamIntegerdataStreamenv.fromElements(1,2,3);TabletableFromStreamtenv.fromDataStream(dataStream);statementSet.add(tableFromStream.insertInto(sinkDescriptor));// attach both pipelines to StreamExecutionEnvironment (the statement set will be cleared after calling this method)statementSet.attachAsDataStream();// define other DataStream API partsenv.fromElements(4,5,6).addSink(newDiscardingSink());// use DataStream API to submit the pipelinesenv.execute();// 1 I[287849559, true]// I[1]// I[2]// I[3]// 3 I[-1058230612, false]// 2 I[-995481497, false]}}以上本文介绍了将table api管道加入datastream。

相关推荐

Substrate Glutton Pallet 深度解析:用 `on_idle` 精准消耗区块权重,把链压榨到极限
Substrate Glutton Pallet 深度解析:用 `on_idle` 精准消耗区块权重,把链压榨到极限

区块链开发框架后端 【免费下载链接】substrate Substrate: The platform for blockchain innovators 项目地址: https://gitcode.com/gh_mirrors/su/substrate 点击查看 免费下载 导读 pallet-glutton(Glutton Pallet)是 Substrate 生态中… · 2026/9/26 2:00:25

基于featexp特征提取方法使用xgboost进行数据分析
基于featexp特征提取方法使用xgboost进行数据分析

在现代金融领域,信贷风险评估是数据科学的一个重要应用场景。通过有效的特征工程和模型优化,可以极大地提高预测的准确性,从而帮助金融机构更好地做出决策。 本文将以特征工程和模型优化为核心,展示如何通过实际案例来解决信贷风险评估中的关键问题。整个过程涵盖了从数据… · 2026/9/26 2:00:25

ZYNQ课设实战:FFT与打地鼠的软硬协同设计指南
ZYNQ课设实战:FFT与打地鼠的软硬协同设计指南

简介:本资源为ZYNQ课设合集,面向电子信息、嵌入式及FPGA方向的学生与开发者,包含「基于ZYNQ的FFT设计与实现」和「基于ZYNQ打地鼠游戏设计」两个完整项目,帮助读者在Xilinx ZYNQ SoC平台上打通软硬件协同开发流程。压缩包为rar格式… · 2026/9/26 2:00:25

用LangFlow搭建流量包推荐智能客服:RAG与对话记忆实战
用LangFlow搭建流量包推荐智能客服:RAG与对话记忆实战

简介:基于LangFlow框架的零代码大模型应用开发平台项目包,面向希望快速搭建智能客服与RAG应用的开发者、产品经理及运维人员。项目以“流量包推荐智能客服”为实战场景,完整演示对话记忆、检索增强生成(RAG)和多种模型… · 2026/9/26 2:37:20

游戏加加监控配置与帧数显示排查全攻略:从原理到实战
游戏加加监控配置与帧数显示排查全攻略:从原理到实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/26 2:37:20

基于迁移学习的昆虫识别系统实战:ResNet50+PyTorch+Flask部署
基于迁移学习的昆虫识别系统实战:ResNet50+PyTorch+Flask部署

简介:基于Python开发的昆虫识别系统,面向毕业设计、课程设计及实际项目开发,提供高精度识别能力与完整工程源码。系统采用模型迭代方式持续优化,最新版已支持2037个昆虫分类单元,Top1/Top5准确率分别达0.922/0.981&… · 2026/9/26 2:37:20

深入解析 zap 官方 FAQ:Go 高性能结构化日志库的设计取舍、采样机制与实战问答(基于 confd 仓库源码验证)
深入解析 zap 官方 FAQ:Go 高性能结构化日志库的设计取舍、采样机制与实战问答(基于 confd 仓库源码验证)

后端配置中心运维 【免费下载链接】confd Manage local application configuration files using templates and data from etcd or consul 项目地址: https://gitcode.com/gh_mirrors/co/confd 点击查看 免费下载 导读 go.uber.org/zap(以下简称 zap&a… · 2026/9/26 2:37:20

ng-zorro-antd 表单动态增减表单项实战:基于 FormArray 实现增删字段、行内布局与校验提交
ng-zorro-antd 表单动态增减表单项实战:基于 FormArray 实现增删字段、行内布局与校验提交

UI组件前端 【免费下载链接】ng-zorro-antd Angular UI Component Library based on Ant Design 项目地址: https://gitcode.com/gh_mirrors/ng/ng-zorro-antd 点击查看 免费下载 动态增加、减少表单项(Dynamic Form Item)是数据录入类页面最… · 2026/9/26 2:37:20

abogen:3步把EPUB变成带同步字幕的有声书
abogen:3步把EPUB变成带同步字幕的有声书

abogen:3步把EPUB变成带同步字幕的有声书 【免费下载链接】abogen Generate audiobooks from EPUBs, PDFs and text with synchronized captions. 项目地址: https://gitcode.com/GitHub_Trending/ab/abogen 想做有声书却不想订录音棚?abogen 走的… · 2026/9/26 2:37:13

数据库课后习题答案别硬背:当测试用例集刷,效率翻倍
数据库课后习题答案别硬背:当测试用例集刷,效率翻倍

简介:万常选版《数据库原理与设计》课后习题答案资源,覆盖第2至6章及第9章,适合正在学习关系模型、数据库建模、关系数据理论与模式求精的本科生、自学者作为复习与自测材料。压缩包共7个文件,含3个doc参考答案、2个sql示例脚本、… · 2026/9/26 0:00:21

OpenClaw 替代品?Hermes Agent 踩坑实录:macOS 飞书接入 TaoToken 配置
OpenClaw 替代品?Hermes Agent 踩坑实录:macOS 飞书接入 TaoToken 配置

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/26 0:00:40

向下兼容与向上兼容:接口设计中的兼容性策略与工程实践
向下兼容与向上兼容:接口设计中的兼容性策略与工程实践

一次版本升级事故,是很多团队绕不过去的坎。线上环境里,服务端明明已经上线了新版接口,老的移动端还在照着旧文档传参数。请求一到网关,校验直接拒绝,用户操作失败,客服群炸了锅,开发群里开始互… · 2026/9/26 0:00:46

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码