跳转至

博客

大家好,我是 xliuqq.


Metrics

示例代码

Dropwizard Metrics

Spark 使用其作为指标框架。

Dropwizard Metrics 能够从各个角度度量已存在的java应用的成熟框架,简便地以jar包的方式集成进系统,可以以http、ganglia、graphite、log4j等方式提供全栈式的监控视野。

Maven

<dependencies>
    <!--https://github.com/dropwizard/dropwizard/releases-->
    <dependency>
        <groupId>io.dropwizard.metrics</groupId>
        <artifactId>metrics-core</artifactId>
        <version>${metrics.version}</version>
    </dependency>

    <dependency>
        <groupId>io.dropwizard.metrics</groupId>
        <artifactId>metrics-healthchecks</artifactId>
        <version>${metrics.version}</version>
    </dependency>
</dependencies>

Registry

度量的核心是MetricRegistry类,它是所有应用程序度量的容器,MetricRegistry是线程安全类

每个registry里的metric都有唯一的名,以 '.' 分隔,如"thing.count",同一个名字对应同一个metric

// 静态方法,生成唯一的名,"key.jobs.size"
String count = MetricRegistry.name("key", "jobs", "size")
// 注册指标
MetricRegistry metrics = new MetricRegistry();
// 1. register 注册
metrics.register(count, (Gauge<Integer>) () -> 5)
// 2. 通过 gauge, meter, counter 等注册
metrics.gauge("key.jobs.size")

指标类型

2023-03-22 17:09:02.016 [INFO ] c.c.m.Slf4jReporter$InfoLoggerProxy.log[473 line]- type=GAUGE, name=call.times, value=10
2023-03-22 17:09:02.016 [INFO ] c.c.m.Slf4jReporter$InfoLoggerProxy.log[473 line]- type=COUNTER, name=count, count=-10
2023-03-22 17:09:02.016 [INFO ] c.c.m.Slf4jReporter$InfoLoggerProxy.log[473 line]- type=HISTOGRAM, name=histogram, count=10, min=0, max=2, mean=0.9639057667940371, stddev=0.7820338303861577, p50=1.0, p75=2.0, p95=2.0, p98=2.0, p99=2.0, p999=2.0
2023-03-22 17:09:02.016 [INFO ] c.c.m.Slf4jReporter$InfoLoggerProxy.log[473 line]- type=METER, name=requests, count=10, m1_rate=0.38646381163615817, m5_rate=0.3968026648037416, m15_rate=0.398904212905539, mean_rate=0.3524120972991338, rate_unit=events/second
Meters

meter测量事件随时间变化的速率,以及1分钟、5分钟、15分钟内的移动平均值

private final MetricRegistry metrics = new MetricRegistry();
private final Meter requests = metrics.meter("requests");

// 测试每分钟请求的频率
public void handleRequest(Request request, Response response) {
    requests.mark();
    // etc
}
Gauges

gauge: 量表是对一个值的瞬时测量。(比如 CPU 的使用率等)

// 只是注册Guage这个metric,Reporter获取的时候才会触发计算(即每次都会进行调用,获取最新的值)
metrics.register(MetricRegistry.name(String.class, "test", "size"),
        (Gauge<Integer>) () -> 5);

默认提供JmxAttributeGaugeRatioGaugeCachedGaugeDerivativeGauge

Counters

可用于统计单次时间,一般用来统计一直增长的指标。

是一个AtomicLong实例的gauge,可以执行increment 和 decrement函数。

// 使用 #counter(String) 而不是 #register(String, Metric)
Counter pendingJobs = metrics.counter(name(String.class, "pending-jobs"));
pendingJobs.inc(); pendingJobs.dec();
Histograms

reservoir sampling:蓄水池采样算法(避免统计所有的数据)

  • UniformReservoir:随机采样
  • SlidingWindowReservoir:只保留最后的 N 个值;

  • SlidingTimeWindowArrayReservoir :滑动窗口采样

  • ExponentiallyDecayingReservoir:默认,指数采样

直方图测量数据流中值的统计分布。除了最小值、最大值、平均值等,它还测量中位数、第75、90、95、98、99和99.9个百分点。

private final Histogram responseSizes = metrics.histogram(name(String.class, "rsizes"));

public void handleRequest(Request request, Response response) {
    // etc
    responseSizes.update(response.getContent().length);
}
Timers

用于统计函数的调用次数和平均执行时间

计时器测量调用特定代码段的速率及其持续时间的分布,包含 Meter 和 Historm。

private final Timer responses = metrics.timer(name(RequestHandler.class, "responses"));
// 计时器将以纳秒为单位测量处理每个请求所需的时间,并提供每秒请求的速率
public String handleRequest(Request request, Response response) {
    try(final Timer.Context context = responses.time()) {
        // etc;
        return "OK";
    } // catch and final logic goes here
}

Health Checks

检查用户服务的健康状态

// 继承HealthCheck类,实现check方法
// 异步执行,每 2s 执行一次
@Async(period = 2000)
public class ServerHealthCheck extends HealthCheck {
    @Override
    public HealthCheck.Result check() throws Exception {
        double value = Math.random();

        if (value > 0.5) {
            return HealthCheck.Result.healthy();
        } else {
            return HealthCheck.Result.unhealthy("error random");
        }
    }
}

// 定义HealthCheckRegistry,默认 2线程的线程池运行异步的check
final HealthCheckRegistry healthChecks = new HealthCheckRegistry();
healthChecks.register("server.healthy", new ServerHealthCheck());
// 运行所有注册的健康检查
final Map<String, HealthCheck.Result> results = healthChecks.runHealthChecks();
for (Entry<String, HealthCheck.Result> entry : results.entrySet()) {
    if (entry.getValue().isHealthy()) {
        System.out.println(entry.getKey() + " is healthy");
    } else {
        System.err.println(entry.getKey() + " is UNHEALTHY: " + entry.getValue().getMessage());
        final Throwable e = entry.getValue().getError();
        if (e != null) {
            e.printStackTrace();
        }
    }
}

Metrics内置ThreadDeadlockHealthCheck健康检查,使用Java内置的线程死锁检测

Reporter

core 内置 ConsoleReporter, CsvReporterSlf4jReporter

Console Reporter

Reporter会定时去MetricRegistry中获取metric结果

ConsoleReporter reporter = ConsoleReporter.forRegistry(metrics)
       .convertRatesTo(TimeUnit.SECONDS)
       .convertDurationsTo(TimeUnit.MILLISECONDS)
       .build();
reporter.start(1, TimeUnit.SECONDS); // 每秒将结果输出到控制台
Reporting Via JMX
<dependency>
    <groupId>io.dropwizard.metrics</groupId>
    <artifactId>metrics-jmx</artifactId>
    <version>${metrics.version}</version>
</dependency>
final JmxReporter reporter = JmxReporter.forRegistry(registry).build();
reporter.start();

一旦reporter开始,所有注册的metrics都可以通过JConsole或者VisualVM查看。

Reporting Via HTTP

AdminServlet 提供 所有注册metrics的Json表示;也会运行健康检查;打印thread dump;对load-balancers提供简单的'ping'返回。

<dependency>
    <groupId>io.dropwizard.metrics</groupId>
    <artifactId>metrics-servlets</artifactId>
    <version>${metrics.version}</version>
</dependency>
Other Reporting

Instrumenting

https://metrics.dropwizard.io/4.2.0/manual/servlet.html

对常见的框架,提供注入,自动记录Metrics,如EhcacheCaffineApache HttpClientJDBILog4jLogBackJettyJersey 2.xJVM

Web-Application的 metrics(metrics-servlet): status codes(meters), the number of active requests(counter), request duration(timer)

  • 通过Filter(com.codahale.metrics.servlet.InstrumentedFilter)

三方库

https://metrics.dropwizard.io/4.2.0/manual/third-party.html

Prometheus Metrics Client

Prometheus instrumentation library for JVM applications

指标类型

Counter

Counters go up, and reset when the process restarts.

示例:

import io.prometheus.client.Counter;
class YourClass {
  static final Counter requests = Counter.build()
     .name("requests_total").help("Total requests.").register();

  void processRequest() {
    requests.inc();
    // Your code here.
  }
}
Gauge

Gauges can go up and down.

示例

class YourClass {
  static final Gauge inprogressRequests = Gauge.build()
     .name("inprogress_requests").help("Inprogress requests.").register();

  void processRequest() {
    inprogressRequests.inc();
    // Your code here.
    inprogressRequests.dec();
  }
}
Summary

monitor distributions, like latencies or request sizes.

默认提供 sum、count 指标,可添加 95%, 90% 的统计信息。

示例

private static final Summary requestLatency = Summary.build()
    .name("requests_latency_seconds")
    .help("request latency in seconds")
    .quantile(0.5, 0.01)    // 0.5 quantile (median) with 0.01 allowed error
    .quantile(0.95, 0.005)  // 0.95 quantile with 0.005 allowed error
    .register();

private static final Summary receivedBytes = Summary.build()
    .name("requests_size_bytes")
    .help("request size in bytes")
    .register();

public void processRequest(Request req) {
    // 计时
    Summary.Timer requestTimer = requestLatency.startTimer();
    try {
        // Your code here.
    } finally {
        requestTimer.observeDuration();
        // 数量统计
        receivedBytes.observe(req.size());
    }
}

时窗:一定时间内的 Summary 而不是整个APP生命周期,可以通过滑动窗口配置

Summary requestLatency = Summary.build()
    .name("requests_latency_seconds")
    .help("Request latency in seconds.")
    .maxAgeSeconds(10 * 60)
    .ageBuckets(5)
    // ...
    .register();

10分钟的时间窗口,5个bucket,即每2分钟滑动一次。

Histogram

monitor distributions, like latencies or request sizes.

每个桶中累积值的数量,bucket 可以配置,形成柱状图。

如果需要计算 φ-quantiles,需要在server端通过histogram_quantile()函数实现。

Labels

所有的指标可以有标签,进行组合。

class YourClass {
  static final Counter requests = Counter.build()
     .name("my_library_requests_total").help("Total requests.")
     .labelNames("method").register();

  void processGetRequest() {
    requests.labels("get").inc();
    // Your code here.
  }
}

指标注册

推荐使用静态变量的方式进行注册(有全局的默认registry)。

static final Counter requests = Counter.build()
   .name("my_library_requests_total").help("Total requests.").labelNames("path").register();

Exemplars

supported for Counter and Histogram

OpenMetrics 格式特性,允许将 metrics 跟 traces 关联。

DefaultExemplarSampler 内置支持 OpenTelemetry tracing

Collectors

默认包含 GC,内存池,类加载和线程数量的指标。

DefaultExports.initialize();
Logging

logging collectors for log4j, log4j2 and logback.

Caches

Guava/Caffeine cache collector

Hibernate
Jetty
Servlet Filter
Spring AOP

add simpleclient_spring_web as a dependency, annotate a configuration class with @EnablePrometheusTiming, then annotate your Spring components as such:

@Controller
public class MyController {
  @RequestMapping("/")
  @PrometheusTimeMethod(name = "my_controller_path_duration_seconds", help = "Some helpful info here")
  public Object handleMain() {
    // Do something
  }
}
DropwizardExports Collector

将 Dropwizard 的指标,转换成 prometheus 的指标。

// Dropwizard MetricRegistry
MetricRegistry metricRegistry = new MetricRegistry();
new DropwizardExports(metricRegistry).register();
  • 将不支持的字符,转换为"_",如
Dropwizard metric name:
org.company.controller.save.status.400
Prometheus metric:
org_company_controller_save_status_400

Exporting

HTTP

client library 中,提供HTTP, Servlet, SpringBoot, Vert.x 集成。

HTTPServer server = new HTTPServer.Builder()
    .withPort(1234)
    .build();

Spark Accumlator 解析

问题:

  • Accumlator 的更新粒度是 Task 级别,还是什么粒度?
  • Task 失败时的累加器信息,是否仍会更新,出现重复的情况?
  • Action 中的累加器如何保证只执行一次?
  • Transform 中的累加器是否有办法解决重算时多次更新的问题?

消息实时推送方案

需求:客户端及时收到服务端的消息通知

  • 对于 Gitee 等网站,其消息都是界面刷新时才通知,注意实时的性价比

短轮询

短轮询:指定的时间间隔,由浏览器向服务器发出HTTP请求,服务器实时返回未读消息数据给客户端,浏览器再做渲染显示。

  • 实现:JS定时器;
  • 缺点:推送数据并不会频繁变更,无论后端此时是否有新的消息产生,客户端都会进行请求,势必会对服务端造成很大压力,浪费带宽和服务器资源。

长轮询(Comet)

示例代码见:LongPullingCode

客户端主动拉pull模型,应用长轮询(Long Polling

  • 由服务端控制响应客户端请求的返回时间,来减少客户端无效请求的一种优化手段;

  • 客户端发起请求后,服务端不会立即返回请求结果,而是将请求挂起等待一段时间,如果此段时间内服务端数据变更,立即响应客户端请求,若是一直无变化则等到指定的超时时间后响应请求,客户端重新发起长链接

长轮询

实现

  • DeferredResultservelet3.0后经过Spring封装提供的一种异步请求机制

  • 允许容器线程快速释放占用的资源,不阻塞请求线程,以此接受更多的请求提升系统的吞吐量。

  • 启动异步工作线程处理真正的业务逻辑,处理完成调用DeferredResult.setResult(200)提交响应结果。

  • 一个ID可能会被多个长轮询请求监听,使用guava包的Multimap结构,一个key对应多个value。一旦监听到key变化,对应的所有长轮询都会响应;

  • 请求超过设置的超时时间,会抛出AsyncRequestTimeoutException异常,全局捕获统一返回,前端获取约定好的状态码后再次发起长轮询请求。

  • 前端得到非请求超时的状态码,知晓数据变更,主动查询未读消息数接口,更新页面数据;

注意

  • 竞态分析:如果在前端重新建立连接时,后端接收到新消息,此时前端还没有建立连接,该消息会丢失,需要根据业务设计相应方案;

缺点

长轮询相比于短轮询在性能上提升了很多,但依然会产生较多的请求,在响应之后,会引起请求突然激增。

SSE

不应该设置 WriteTimeOut,示例代码见:SSE 代码

服务器发送事件(Server-sent events),简称SSE

SSE在服务器和客户端之间打开一个单向通道,服务端响应的不再是一次性的数据包而text/event-stream类型的数据流信息,在有数据变更时从服务器流式传输到客户端。

注:SSE不支持IE浏览器,对其他主流浏览器兼容性做的还不错。

SSE与Websocket的区别

WebSocket SSE
全双工,可以同时发送和接收消息 单工,只能服务端单向发送消息
独立的协议 基于HTTP协议
需要代理服务器单独支持 代理服务器直接支持
协议相对复杂 协议相对简单,易于理解和使用
默认不支持断线重连 默认支持断线重连
默认支持传送二进制数据 一般只用来传送文本,二进制数据需要编码后传送
不支持自定义发送的数据类型 支持自定义发送的数据类型
支持CORS 不支持CORS,协议和端口都必须相同
  • 以1次/秒或者更快的频率向服务端传输数据,那应该用WebSocket;
  • 客户端和服务端脚本之间具有网络服务器情况时,一个SSE连接不仅使用一个套接字,还会占用一个Apache线程或进程;

实现

客户端:通过 EventSource建立连接;

  • 会自动断线重连(如超时重连),并且会发送Last-Event-ID header

服务端:Springboot 通过 SseEmitter实现;

  • SSE消息的格式为field:value\nfield:value\n\n, field 支持如下
: 即以:开头表示注释可以理解为服务端向客户端发送的心跳确保连接不中断
data数据
event: 事件默认值为 message如果是其它事件前端可以通过 addEventListener 监听不同的事件
id: 数据标识符用 id 字段表示相当于每一条数据的编号
retry: 重连时间

Nginx 转发 SSE

  • nginx upstream 超时连接关闭(upstream timeout)时,前端会自动进行重连;
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

MQTT 全称(Message Queue Telemetry Transport):一种基于发布/订阅(publish/subscribe)模式的轻量级通讯协议,通过订阅相应的主题来获取消息,是物联网Internet of Thing)中的一个标准传输协议。

将消息的发布者(publisher)与订阅者(subscriber)进行分离,因此可以在不可靠的网络环境中,为远程连接的设备提供可靠的消息服务,使用方式与传统的MQ有点类似。

实现

https://mp.weixin.qq.com/s/U-fUGr9i1MVa4PoVyiDFCg

WebSocket

TCP连接上进行全双工通信的协议

SOCKJS

当今主流浏览器都支持 Websocket,因此 SOCKJS 不再需要

  • 除非反向代理(如Nginx)的配置不支持WebSocket.

使用 SOCKJS (模拟 Websocket 协议)

  • sockjs-client 提供浏览器兼容性,优先使用原生的WebSocket,如果某个浏览器不支持WebSocket,SockJS会自动降级为SSE 或者 长轮询;

心跳

  • 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

使用 STOMP(定义Websocket的数据格式),消息的订阅和发送。

  • STOMP即Simple Text Orientated Messaging Protocol,简单文本定向消息协议;
  • stompjs 客户端,1.1 版本提供心跳机制(默认10s,但会跟 broker 连接时协调出最终的in/out心跳);
  • 连接稳定性,支持自动重连;

事务支持

// 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();

Spring

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;

  • 跨域问题,the default behavior for WebSocket and SockJS is to accept only same-origin requests. ;
  • nginx 配置 proxy_set_header Host $host;
  • 或者 spring 配置 setAllowedOriginPatterns("*")

  • 事件问题:

  • SessionDisconnectEvent需要幂等处理,一个WS连接(Session)可能会被发送多次;
Simple Broker

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):

外部 Broker

支持 RabbitMQ, ActiveMQ 等

认证

通过configureClientInboundChannel()注册消息处理,进行验证。

nginx 转发 Websocket

proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_http_version 1.1;

proxied server 端需要定时发送 PING 消息保证 nginx 不会关闭连接。

站内信设计

实现一个站内信web消息推送的功能

img

消息推送(push) 通常是指网站的运营工作等人员,通过某种工具对用户当前网页或移动设备APP进行的主动消息推送。

消息推送一般又分为web端消息推送移动端消息推送

web端消息推送常见的诸如站内信、未读邮件数量、监控报警数量等,应用的也非常广泛。

需求:站内信

  • 管理员可以发布多种不同的消息,用户在界面可以及时收到通知;

  • 每个 client,需要有自己的唯一标识,防止用户对一个界面同时多开,此时每个界面都会有一个独立的连接;

  • 即仅用用户ID作为 ClientId 不行,因为 Server 端对同一个ID,只会保留一个连接,但是对应多个前端页面连接,会出错;

非实时方案

类似 Gitee,其站内信是非实时的,需要界面刷新才可以获取。

  • unread接口,获取未读消息;

实时方案

触发某个事件(主动分享了资源或者后台主动推送消息),web页面的通知小红点就会实时的+1

在服务端会有若干张消息推送表,用来记录用户触发不同事件所推送不同类型的消息,前端主动查询(拉)或者被动接收(推)用户所有未读的消息数

接口管理

接口版本间改动

微服务的不同版本间,可以查看接口的改动。

基本思路:每个版本都有对应的 OpenAPI 规范的 Yaml 描述文件,通过OpenAPI Tools进行接口版本比对;

每个版本如何具备对应的OpenAPI的yaml描述文件?

思路1:根据Yaml生成Controller层注解代码

开发接口先写OpenAPI 3 的 yaml 文件,然后生成对应的Controller层的接口。(可行,但对于开发不太友好,注解比yaml容易写)

<generateApiDocumentation>false</generateApiDocumentation>
<generateSupportingFiles>false</generateSupportingFiles>
<skipOperationExample>true</skipOperationExample>

具体代码见 根据Yaml生成SpringBoot接口

思路2:根据Controller层注解代码生成Yaml

github上有直接根据注解进行静态解析,生成对应的 yaml 接口文档,但是鲁棒性不够,无生产可用;

工具:

  • springdoc-openapi : 在SpringBoot中使用OpenAPI 3的注解;

  • springdoc-openapi-maven-plugin : 根据 OpenAPI 3的注解生成OpenAPI 3的Yaml,依赖于 spring-boot-maven 插件

  • Maven在集成测试阶段(integration-test)运行 openapi 插件(通过mvn verify)。

问题点:

  • 但是在CI阶段,一般不进行verify,仅package不会运行spring,也就无法生成 Yaml;

具体代码见 根据注解生成Yaml

接口文档注解生成

OpenAPI/Swagger 3规范

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;

OpenAPI2.0 OpenAPI3.0 info

springboot 集成(springdoc

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")

示例项目见

基于Tag自动发布Github Release

在 Github 上发布 Release 时,默认只会发布源码的zip包,但是有时候会想要将一些编译出来的文件也放在 Release 中。虽然可以在本地编译,然后通过界面手动上传。但是可以基于 Github Action 实现自动发布

Maven 插件开发

在开发通过包/类/方法注解自动生成README的功能时,在使用时,如果使用 Jar 依赖的方式,需要为每个工程定义一个 Main 方法,并且需要手动调用。

因此想通过开发Maven 插件,在compile阶段自动调用并生成 README,因此便研究如何开发 Maven 插件。

Spark 使用 Yarn Rest 提交

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

K8s 中 leader election 选举原理

在开发CRD时,定义 controller 的时候,会看到如下代码

mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
    // 是否进行 leader 选举
    LeaderElection:          enableLeaderElection,
    // Namespace and name
    LeaderElectionNamespace: leaderElectionNamespace,
    LeaderElectionID:        "alluxio.data.fluid.io",
    // ...
})

对于有状态组件来说,实现高可用一般来说通过选主来达到同一时刻只能有一个组件在处理业务逻辑。

这里,会比较好奇这个选举是如何实现的,接下来的内容便从源码的角度进行解读。

Maven 自定义仓库

当自己开发一个工具包,然后另一个项目要引用时,因此需要将 jar 包放到可访问的公网上:

  • Maven 官方 repo,如使用 Sonatype OSSRH,但其注册复杂,此次不进行介绍;

  • 第三方公共仓库,如 JitPack, Github/Gitee等,后面进行介绍。

K8s Pod 如何配置免密

MPI Operator 在执行 MPI 作业时,通过 SSH 的方式,因此需要在不同的 MPI Workers Pod间配置免密。

本文通过源码分析,探究如何在 K8s Pod 间配置免密。