Spark Accumlator 解析
问题:
- Accumlator 的更新粒度是 Task 级别,还是什么粒度?
- Task 失败时的累加器信息,是否仍会更新,出现重复的情况?
- Action 中的累加器如何保证只执行一次?
- Transform 中的累加器是否有办法解决重算时多次更新的问题?
问题:
需求:客户端及时收到服务端的消息通知
- 对于 Gitee 等网站,其消息都是界面刷新时才通知,注意实时的性价比
短轮询:指定的时间间隔,由浏览器向服务器发出HTTP请求,服务器实时返回未读消息数据给客户端,浏览器再做渲染显示。
示例代码见:LongPullingCode
客户端主动拉pull模型,应用长轮询(Long Polling)
由服务端控制响应客户端请求的返回时间,来减少客户端无效请求的一种优化手段;
客户端发起请求后,服务端不会立即返回请求结果,而是将请求挂起等待一段时间,如果此段时间内服务端数据变更,立即响应客户端请求,若是一直无变化则等到指定的超时时间后响应请求,客户端重新发起长链接;

DeferredResult 在servelet3.0后经过Spring封装提供的一种异步请求机制
允许容器线程快速释放占用的资源,不阻塞请求线程,以此接受更多的请求提升系统的吞吐量。
启动异步工作线程处理真正的业务逻辑,处理完成调用DeferredResult.setResult(200)提交响应结果。
一个ID可能会被多个长轮询请求监听,使用guava包的Multimap结构,一个key对应多个value。一旦监听到key变化,对应的所有长轮询都会响应;
请求超过设置的超时时间,会抛出AsyncRequestTimeoutException异常,全局捕获统一返回,前端获取约定好的状态码后再次发起长轮询请求。
前端得到非请求超时的状态码,知晓数据变更,主动查询未读消息数接口,更新页面数据;
注意:
长轮询相比于短轮询在性能上提升了很多,但依然会产生较多的请求,在响应之后,会引起请求突然激增。
不应该设置 WriteTimeOut,示例代码见:SSE 代码
服务器发送事件(Server-sent events),简称SSE。
SSE在服务器和客户端之间打开一个单向通道,服务端响应的不再是一次性的数据包而是text/event-stream类型的数据流信息,在有数据变更时从服务器流式传输到客户端。
注:SSE不支持IE浏览器,对其他主流浏览器兼容性做的还不错。
| WebSocket | SSE |
|---|---|
| 全双工,可以同时发送和接收消息 | 单工,只能服务端单向发送消息 |
| 独立的协议 | 基于HTTP协议 |
| 需要代理服务器单独支持 | 代理服务器直接支持 |
| 协议相对复杂 | 协议相对简单,易于理解和使用 |
| 默认不支持断线重连 | 默认支持断线重连 |
| 默认支持传送二进制数据 | 一般只用来传送文本,二进制数据需要编码后传送 |
| 不支持自定义发送的数据类型 | 支持自定义发送的数据类型 |
| 支持CORS | 不支持CORS,协议和端口都必须相同 |
客户端:通过 EventSource建立连接;
Last-Event-ID header服务端:Springboot 通过 SseEmitter实现;
field:value\nfield:value\n\n, field 支持如下空: 即以:开头,表示注释,可以理解为服务端向客户端发送的心跳,确保连接不中断
data:数据
event: 事件,默认值为 message,如果是其它事件,前端可以通过 addEventListener 监听不同的事件
id: 数据标识符用 id 字段表示,相当于每一条数据的编号
retry: 重连时间
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_http_version 1.1;
# Nginx 默认 OFF
proxy_cache off;
# 应该由 server 端发送 X-Accel-Buffering: no;
# 针对 proxy_pass 的后端配置,好像不需要配置,即使开启前端也会实时响应
# proxy_buffering off;
# 该配置对于SSE不需要
# chunked_transfer_encoding off;
MQTT 全称(Message Queue Telemetry Transport):一种基于发布/订阅(publish/subscribe)模式的轻量级通讯协议,通过订阅相应的主题来获取消息,是物联网(Internet of Thing)中的一个标准传输协议。
将消息的发布者(publisher)与订阅者(subscriber)进行分离,因此可以在不可靠的网络环境中,为远程连接的设备提供可靠的消息服务,使用方式与传统的MQ有点类似。
https://mp.weixin.qq.com/s/U-fUGr9i1MVa4PoVyiDFCg
在TCP连接上进行全双工通信的协议
当今主流浏览器都支持 Websocket,因此 SOCKJS 不再需要
- 除非反向代理(如Nginx)的配置不支持WebSocket.
使用 SOCKJS (模拟 Websocket 协议)
心跳
SockJS 协议要求 servers 发送心跳消息,排除代理导致的连接挂起。(Spring SockJS 默认 heartbeatTime 为 25s);
When using STOMP over WebSocket and SockJS, if the STOMP client and server negotiate heartbeats to be exchanged, the SockJS heartbeats are disabled.
使用 STOMP(定义Websocket的数据格式),消息的订阅和发送。
事务支持
// start the transaction
const tx = client.begin();
// send the message in a transaction
client.publish({
destination: '/queue/test',
headers: { transaction: tx.id },
body: 'message in a transaction',
});
// commit the transaction to effectively send the message
tx.commit();
https://docs.spring.io/spring-framework/docs/current/reference/html/web.html#websocket-server
On the Servlet stack, the Spring Framework provides both server (and also client) support for the SockJS protocol;
proxy_set_header Host $host;或者 spring 配置 setAllowedOriginPatterns("*");
事件问题:
SessionDisconnectEvent需要幂等处理,一个WS连接(Session)可能会被发送多次;The built-in simple message broker handles subscription requests from clients, stores them in memory, and broadcasts messages to connected clients that have matching destinations.
纯 WebSocket 的心跳问题(不使用 SockJS):
通过configureClientInboundChannel()注册消息处理,进行验证。
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_http_version 1.1;
proxied server 端需要定时发送 PING 消息保证 nginx 不会关闭连接。
实现一个站内信web消息推送的功能

消息推送(push) 通常是指网站的运营工作等人员,通过某种工具对用户当前网页或移动设备APP进行的主动消息推送。
消息推送一般又分为web端消息推送和移动端消息推送。
web端消息推送常见的诸如站内信、未读邮件数量、监控报警数量等,应用的也非常广泛。
需求:站内信
管理员可以发布多种不同的消息,用户在界面可以及时收到通知;
每个 client,需要有自己的唯一标识,防止用户对一个界面同时多开,此时每个界面都会有一个独立的连接;
类似 Gitee,其站内信是非实时的,需要界面刷新才可以获取。
unread接口,获取未读消息;触发某个事件(主动分享了资源或者后台主动推送消息),web页面的通知小红点就会实时的+1。
在服务端会有若干张消息推送表,用来记录用户触发不同事件所推送不同类型的消息,前端主动查询(拉)或者被动接收(推)用户所有未读的消息数。
微服务的不同版本间,可以查看接口的改动。
基本思路:每个版本都有对应的 OpenAPI 规范的 Yaml 描述文件,通过OpenAPI Tools进行接口版本比对;
每个版本如何具备对应的OpenAPI的yaml描述文件?
开发接口先写OpenAPI 3 的 yaml 文件,然后生成对应的Controller层的接口。(可行,但对于开发不太友好,注解比yaml容易写)
<generateApiDocumentation>false</generateApiDocumentation>
<generateSupportingFiles>false</generateSupportingFiles>
<skipOperationExample>true</skipOperationExample>
具体代码见 根据Yaml生成SpringBoot接口
github上有直接根据注解进行静态解析,生成对应的 yaml 接口文档,但是鲁棒性不够,无生产可用;
工具:
springdoc-openapi : 在SpringBoot中使用OpenAPI 3的注解;
springdoc-openapi-maven-plugin : 根据 OpenAPI 3的注解生成OpenAPI 3的Yaml,依赖于 spring-boot-maven 插件
openapi 插件(通过mvn verify)。问题点:
具体代码见 根据注解生成Yaml
https://github.com/swagger-api/swagger-codegen
支持json和yaml文件解析,自动生成多种语言的API客户端和服务器stub;
swagger配置规范说明,https://swagger.io/specification;
swagger maven plugin,https://github.com/garethjevans/swagger-codegen-maven-plugin;

springfox 不支持 springboot 2.4 以上的版本,最新更新时间为 2020/10/14号。
springdoc 与 springboot 的版本兼容性:
| spring-boot Versions | Minimum springdoc-openapi Versions |
|---|---|
3.0.x |
2.0.x+ |
2.7.x, 1.5.x |
1.6.0+ |
2.6.x, 1.5.x |
1.6.0+ |
2.5.x, 1.5.x |
1.5.9+ |
2.4.x, 1.5.x |
1.5.0+ |
2.3.x, 1.5.x |
1.4.0+ |
2.2.x, 1.5.x |
1.2.1+ |
2.0.x, 1.5.x |
1.0.0+ |
SpringFox - SpringDoc 对应关系如下
@Api → @Tag@ApiIgnore → @Parameter(hidden = true) or @Operation(hidden = true) or @Hidden@ApiImplicitParam → @Parameter@ApiImplicitParams → @Parameters@ApiModel → @Schema@ApiModelProperty(hidden = true) → @Schema(accessMode = READ_ONLY)@ApiModelProperty → @Schema@ApiOperation(value = "foo", notes = "bar") → @Operation(summary = "foo", description = "bar")@ApiParam → @Parameter@ApiResponse(code = 404, message = "foo") → @ApiResponse(responseCode = "404", description = "foo")示例项目见
在 Github 上发布 Release 时,默认只会发布源码的zip包,但是有时候会想要将一些编译出来的文件也放在 Release 中。虽然可以在本地编译,然后通过界面手动上传。但是可以基于 Github Action 实现自动发布。
在开发通过包/类/方法注解自动生成README的功能时,在使用时,如果使用 Jar 依赖的方式,需要为每个工程定义一个 Main 方法,并且需要手动调用。
因此想通过开发Maven 插件,在compile阶段自动调用并生成 README,因此便研究如何开发 Maven 插件。
https://community.cloudera.com/t5/Community-Articles/Starting-Spark-jobs-directly-via-YARN-REST-API/ta-p/245998 这篇文章是Spark 1.6, 已经过时,只能作为基本的参考,具体还是要阅读Spark on Yarn的提交代码。
- 本文基于 Spark 2.4.8 和 Hadoop 3.2 进行验证。
将 Spark 作业提交到 Yarn上时,只能通过命令行 spark-submit 进行操作,本文通过解析 spark-submit 的源码,探究如何使用 Yarn Rest API 进行提交 Spark 作业(仅 cluster 模式,因 client 模式 driver 运行在 client 中而不是 AM 中)。
一句话总结:还是用命令行调用spark-submit !
在开发CRD时,定义 controller 的时候,会看到如下代码
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
// 是否进行 leader 选举
LeaderElection: enableLeaderElection,
// Namespace and name
LeaderElectionNamespace: leaderElectionNamespace,
LeaderElectionID: "alluxio.data.fluid.io",
// ...
})
对于有状态组件来说,实现高可用一般来说通过选主来达到同一时刻只能有一个组件在处理业务逻辑。
这里,会比较好奇这个选举是如何实现的,接下来的内容便从源码的角度进行解读。