当前位置: 首页 > news >正文

《零基础入门Spark》学习笔记 Day 13

Structured Streaming

数据加载

SparkSession的readStream API 来创建DataFrame

var df: DataFrame = spark.readStream .format("socket") .option("host",host) .option("port",port) .load()

format:指定流处理的数据源头类型

option:与数据源头有关的若干选项

load:将数据流加载进Spark

流计算有3个重要的基础概念,比如flink也是如此

Source:流计算的数据源头

Processing:负责对数据流进行转换、过滤、聚合等操作

Sink:指的是数据流向的目的地

数据处理

/** 使用DataFrame API完成Word Count计算 */ // 首先把接收到的字符串,以空格为分隔符做拆分,得到单词数组words df = df.withColumn("words", split($"value", " ")) // 把数组words展平为单词word .withColumn("word", explode($"words")) // 以单词word为Key做分组 .groupBy("word") // 分组计数 .count()

数据输出

/** 将Word Count结果写入到终端(Console) */ df.writeStream // 指定Sink为终端(Console) .format("console") // 指定输出选项 .option("truncate", false) // 指定输出模式 .outputMode("complete") //.outputMode("update") // 启动流处理应用 .start() // 等待中断指令 .awaitTermination()

一般来说,Structured Streaming支持3种Sink输出模式

Complete mode:输出到目前为止处理过的全部内容

Append mode:仅输出最近一次作业的计算结果

Update mode:仅输出内容有根据输入的计算结果

http://www.jsqmd.com/news/615314/

相关文章:

  • 高效光伏电池建模技术分享:Boost Buck电路实现最大功率追踪
  • Java 线程、进程、CPU缓存、MESI
  • Untrunc视频修复工具:让损坏的MP4文件重获新生
  • 微软常用运行库 安装教程:一键修复VC++运行环境(AIO合集)
  • 将盾CDN:业务安全与反欺诈的实战策略
  • 用 AI Coding 工具生成 万字奇幻世界设定的实践记录滥
  • TrailBase 与 PocketBase 详细对比
  • 将盾CDN:移动应用安全合规的实践指南
  • ESP32/8266利用闪存文件系统创建 Web服务实现交互控制
  • GLM-. 全面支持与 Gemini CLI 集成:HagiCode 的多模型进化之路赂
  • IOFILE结构体的介绍与House of orange轮
  • 深度解析DHCP协议:工作原理、4步交互流程及应用场景
  • std::time
  • mysql如何在本地开发环境模拟生产环境_利用Docker克隆
  • 20个核心AI概念拆解:小白也能轻松入门大模型,收藏这份学习秘籍!
  • vue el-table 切换页面、组件销毁会内存泄漏吗?99% 的人都误解了
  • claude (二) skill
  • AI Agent Harness Engineering 如何改变咨询行业并重新定价知识
  • 详细解析Spring如何解决循环依赖问题事
  • 将盾CDN:云原生环境的安全防护策略
  • gopher-os硬件检测系统:从Multiboot到ACPI的完整硬件抽象层设计
  • 营销自动化数据驱动 - 多源数据 OLAP 架构演进厝
  • OCAD应用:打入式断续变焦光学系统初始结构设计
  • 视频分析神器video-analyzer:5分钟学会AI智能视频内容理解终极指南
  • 2026.4.9
  • OpenClaw多模态扩展:Qwen3-14B镜像对接OCR截图识别
  • 支付密钥硬编码、调试模式未关闭、日志泄露token——PHP生产环境支付接口的3大“自杀式配置”
  • Source Han Serif CN 开源字体:技术特性与多场景适配指南
  • 3大技术突破重新定义多模态交互:AudioCLIP的跨模态语义对齐解决方案
  • LANs.py源码深度剖析:理解多线程异步数据包处理机制