logo
publist
写文章

简介

该用户还未填写简介

擅长的技术栈

可提供的服务

暂无可提供的服务

Flink知识点(五)|Window(窗口)

Flink窗口将无界流切分为有界数据块以进行聚合计算。主要分为滚动窗口(固定大小、不重叠)、滑动窗口(固定大小、有重叠)和会话窗口(按数据间隔动态关闭)。窗口函数包括增量聚合(如ReduceFunction)和全量处理(如ProcessWindowFunction)。Flink SQL通过TVF(如TUMBLE、HOP、SESSION)支持窗口操作,并可结合侧输出流处理迟到数据,是流处理实时统计的

文章图片
#flink#大数据
Flink知识点(三)|Flink中的广播变量(Broadcast State)

Flink广播状态(Broadcast State)是一种用于在流处理中将控制流(如规则、配置)广播到所有并行算子实例的机制。其核心设计是:广播流可写入状态,数据流(主流)只读状态,以此保证分布式环境下状态一致性。典型应用于动态规则下发(如实时风控)或小维表与大流Join等场景。使用时需注意广播状态数据量不宜过大,且需开启Checkpoint以保证容错。

文章图片
#flink#大数据
Flin知识点(六)|Flink状态管理

本文系统剖析了Flink状态管理,涵盖Operator State与Keyed State两大类及其具体实现(如ValueState、MapState等),并详细对比了HashMap与RocksDB两种状态后端的原理、选型与配置。作为流计算容错基石,状态管理支撑着广播变量、窗口聚合等核心功能。

文章图片
#flink#大数据
Flink知识点(一)|Flink中的双流关联

本文档全面梳理了Flink实现双流关联的核心方法。DataStream API提供Window Join、Interval Join和可高度自定义的CoProcessFunction。Flink SQL则支持包括Regular Join、Interval Join、Temporal Join、Lookup Join在内的多种语义。文档通过对比各类关联方式的特点、限制与适用场景(如订单-支付流关联)

文章图片
#flink
Flink知识点(二)|Flink中是怎么处理乱序数据的

Flink 处理乱序数据的核心是采用 **事件时间(Event Time)** 语义,通过 **Watermark** 机制来推断事件时间进度并容忍一定范围的乱序。基于此,**窗口** 在 Watermark 越过窗口结束时触发计算,并通过 **Allowed Lateness** 允许延迟关闭以处理部分迟到数据。最终无法处理的迟到事件可通过 **侧输出(Side Output)** 收集,用于补

文章图片
#flink#大数据
Flink知识点(四)|Watermark(水位线)

Watermark(水位线)是Flink处理基于事件时间(Event Time)乱序数据的核心机制。它通过设置一个延迟阈值(如允许乱序10秒),来判定某个时间点之前的数据是否已全部到达,从而准确触发窗口计算。支持在数据源或流处理中周期性或逐条生成,是保证乱序场景下计算正确性的关键组件。

文章图片
#flink#大数据
Flink知识点(七)|容错机制

本文介绍了Flink的容错机制与生产级调优,重点解析了Checkpoint和Savepoint两种状态快照机制。Checkpoint作为轻量级自动周期性快照,基于Chandy-Lamport算法实现Exactly-Once语义;Savepoint则是手动触发的稳定格式快照,用于版本升级和迁移。文章详细说明了两种机制的配置方法、操作流程及生产实践建议,包括目录管理、状态恢复、并行度调整等关键问题排查

文章图片
#flink#大数据
AI知识点(一)|准备工作-Anaconda使用文档

Anaconda是一款集成了Python和常用科学计算工具(如Jupyter Notebook、NumPy、Pandas)的发行版,自带conda环境管理工具。本文档介绍了Anaconda的下载安装方法(Windows/macOS)、基础使用(包括虚拟环境创建、包管理和代码运行方式)以及常见问题解决方案。特别提供了在PyCharm中配置Conda环境的详细步骤,包括解释器设置、包管理和环境切换等操

文章图片
#python#conda
到底了