使用Apache Beam执行内存处理
在使用 Apache Beam 处理内存中数据时,需要意识到 Beam 并不是一个数据处理引擎,而是一个允许在不同引擎(如 Spark、Flink、Google Dataflow)上创建和运行统一管道的 SDK。为了在内存中处理消息,需要使用受支持的数据处理引擎或 DirectRunner,后者仅用于测试目的且存在限制。
我正在运行自己的 GRPC 服务器,收集来自各种数据源的事件。服务器是用 Go 开发的,所有事件源都以预定义格式作为 protobuf 消息发送事件。
我想要做的是在内存中使用 Apache Beam 处理所有这些事件。
我浏览了 Apache Beam 的文档,但找不到可以实现我想要的功能的示例。我不会使用 Kafka、Flink 或任何其他流媒体平台,只是在内存中处理消息并输出结果。
有人可以告诉我开始编写简单的流处理应用程序的正确方法吗?
解决方案
好吧,首先,Apache Beam 不是一个数据处理引擎,它是一个 SDK,允许您创建统一的管道并在不同的引擎上运行它,例如 Spark、Flink、Google Dataflow 等。所以,要运行Beam 管道,您需要利用任何受支持的数据处理引擎或使用 DirectRunner,它将在本地运行您的管道,但(!)它有很多限制,并且主要是为了测试目的而开发的。
正如 Beam 中的每个管道一样,都必须有一个源转换(有界或无界),它将从数据源读取数据。我可以猜测,在您的情况下,您的 GRPC 服务器应该重新传输收集到的事件。因此,对于源转换,您可以使用已实现的 Beam IO transforms(IO 连接器)或创建您自己的,因为 Beam 中目前没有 GrpcIO 或类似的东西。
关于内存中数据的处理,我不确定我是否完全理解你的意思。它主要取决于所使用的数据处理引擎,因为最终,您的 Beam 管道将在实际运行之前转换为 Spark 或 Flink 管道(如果您相应地使用 SparkRunner 或 FlinkRunner ),然后数据处理引擎将管理管道工作流程。大多数现代引擎尽最大努力将所有已处理的数据保留在内存中,并仅在最后手段下将其刷新到磁盘上。
以上就是《使用Apache Beam执行内存处理》的详细内容,更多关于的资料请关注golang学习网公众号!
使用 GORM 的 Raw() 方法执行查询
- 上一篇
- 使用 GORM 的 Raw() 方法执行查询
- 下一篇
- 简化 go mod init 的命令
-
- Golang · Go问答 | 2天前 |
- Go ColumnTypes 推断动态查询字段的安全用法
- 319浏览 收藏
-
- Golang · Go问答 | 2天前 | Go问答 · Go 错误检查 database/sql Rows.Err
- Go Rows.Err 在迭代结束后的错误检查
- 266浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go sql.Scan NULL 到字符串失败的字段适配
- 388浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go Tx.Stmt 复用预处理语句的生命周期
- 154浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go 事务隔离级别设置未生效的驱动边界
- 165浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go DB PingContext 启动探活的超时配置
- 438浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go database/sql Rows 未关闭造成连接耗尽的诊断
- 354浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go database/sql 连接池 MaxIdleConns 的容量关系
- 478浏览 收藏
-
- Golang · Go问答 | 2天前 |
- Go URL 查询值乱码时的编码排查步骤
- 285浏览 收藏
-
- Golang · Go问答 | 2天前 | Go问答 · Go URL路径 url.JoinPath PathEscape 双斜杠
- Go url.JoinPath 处理双斜杠的路径规则
- 446浏览 收藏
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 543次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 516次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 500次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 485次学习
-
- PubMedQA
- 深入了解PubMedQA生物医学问答数据集,涵盖其核心功能、使用方法及在临床决策、药物研发等场景的应用,助力提升NLP模型性能。
- 284次使用
-
- H2O EvalGPT
- H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
- 339次使用
-
- LMArena
- LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
- 336次使用
-
- HELM
- 深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
- 303次使用
-
- MMBench
- MMBench是由上海人工智能实验室等机构联合推出的多模态基准测试平台,提供细粒度能力评估、大规模数据集及VLMEvalKit工具。本文详细介绍其核心功能、安装使用方法及应用场景,助力开发者全面评估多模态模型性能。
- 124次使用
-
- 用Nginx反向代理部署go写的网站。
- 2023-01-17 502浏览
-
- GoLand调式动态执行代码
- 2023-01-13 502浏览
-
- Go tls.GetCertificate 为什么收不到空 ServerName 请求
- 2026-09-27 501浏览
-
- Go sql.Tx提交成功前读取结果导致事务边界混乱的修复方法
- 2026-09-20 501浏览
-
- Go select 用 time.After 做超时有什么资源代价
- 2026-09-10 501浏览
