Http长轮询数据同步源码分析
Apache ShenYu 是一个异步的,高性能的,跨语言的,响应式的
API网关。
在ShenYu网关中,数据同步是指,当在后台管理系统中,数据发送了更新后,如何将更新的数据同步到网关中。Apache ShenYu 网关当前支持ZooKeeper、WebSocket、Http长轮询、Nacos 、Etcd 和 Consul 进行数据同步。本文的主要内容是基于Http长轮询的数据同步源码分析。
本文基于
shenyu-2.5.0版本进行源码分析,官网的介绍请参考 数据同步原理 。
1. Http长轮询
这里直接引用官网的相关描述:
Zookeeper和WebSocket数据同步的机制比较简单,而Http长轮询则比较复杂。Apache ShenYu借鉴了Apollo、Nacos的设计思想,取其精华,自己实现了Http长轮询数据同步功能。注意,这里并非传统的ajax长轮询!

Http长轮询 机制如上所示,Apache ShenYu网关主动请求 shenyu-admin 的配置服务,读取超时时间为 90s,意味着网关层请求配置服务最多会等待 90s,这样便于 shenyu-admin 配置服务及时响应变更数据,从而实现准实时推送。
Http长轮询 机制是由网关主动请求 shenyu-admin ,所以这次的源码分析,我们从网关这一侧开始。
2. 网关数据同步
2.1 加载配置
Http长轮询 数据同步配置的加载是通过spring boot的starter机制,当我们引入相关依赖和在配置文件中有如下配置时,就会加载。
在pom文件中引入依赖:
<!--shenyu data sync start use http-->
<dependency>
<groupId>org.apache.shenyu</groupId>
<artifactId>shenyu-spring-boot-starter-sync-data-http</artifactId>
<version>${project.version}</version>
</dependency>
在application.yml配置文件中添加配置:
shenyu:
sync:
http:
url : http://localhost:9095
当网关启动时,配置类HttpSyncDataConfiguration就会执行,加载相应的Bean。
/**
* Http sync data configuration for spring boot.
*/
@Configuration
@ConditionalOnClass(HttpSyncDataService.class)
@ConditionalOnProperty(prefix = "shenyu.sync.http", name = "url")
@EnableConfigurationProperties(value = HttpConfig.class)
public class HttpSyncDataConfiguration {
private static final Logger LOGGER = LoggerFactory.getLogger(HttpSyncDataConfiguration.class);
/**
* Rest template.
* 创建RestTemplate
* @param httpConfig the http config http配置
* @return the rest template
*/
@Bean
public RestTemplate restTemplate(final HttpConfig httpConfig) {
OkHttp3ClientHttpRequestFactory factory = new OkHttp3ClientHttpRequestFactory();
factory.setConnectTimeout(Objects.isNull(httpConfig.getConnectionTimeout()) ? (int) HttpConstants.CLIENT_POLLING_CONNECT_TIMEOUT : httpConfig.getConnectionTimeout());
factory.setReadTimeout(Objects.isNull(httpConfig.getReadTimeout()) ? (int) HttpConstants.CLIENT_POLLING_READ_TIMEOUT : httpConfig.getReadTimeout());
factory.setWriteTimeout(Objects.isNull(httpConfig.getWriteTimeout()) ? (int) HttpConstants.CLIENT_POLLING_WRITE_TIMEOUT : httpConfig.getWriteTimeout());
return new RestTemplate(factory);
}
/**
* AccessTokenManager.
* 创建AccessTokenManager,专门用户对admin进行http请求时access token的处理
* @param httpConfig the http config.
* @param restTemplate the rest template.
* @return the access token manager.
*/
@Bean
public AccessTokenManager accessTokenManager(final HttpConfig httpConfig, final RestTemplate restTemplate) {
return new AccessTokenManager(restTemplate, httpConfig);
}
/**
* Http sync data service.
* 创建 HttpSyncDataService
* @param httpConfig the http config
* @param pluginSubscriber the plugin subscriber
* @param restTemplate the rest template
* @param metaSubscribers the meta subscribers
* @param authSubscribers the auth subscribers
* @param accessTokenManager the access token manager
* @return the sync data service
*/
@Bean
public SyncDataService httpSyncDataService(final ObjectProvider<HttpConfig> httpConfig,
final ObjectProvider<PluginDataSubscriber> pluginSubscriber,
final ObjectProvider<RestTemplate> restTemplate,
final ObjectProvider<List<MetaDataSubscriber>> metaSubscribers,
final ObjectProvider<List<AuthDataSubscriber>> authSubscribers,
final ObjectProvider<AccessTokenManager> accessTokenManager) {
LOGGER.info("you use http long pull sync shenyu data");
return new HttpSyncDataService(
Objects.requireNonNull(httpConfig.getIfAvailable()),
Objects.requireNonNull(pluginSubscriber.getIfAvailable()),
Objects.requireNonNull(restTemplate.getIfAvailable()),
metaSubscribers.getIfAvailable(Collections::emptyList),
authSubscribers.getIfAvailable(Collections::emptyList),
Objects.requireNonNull(accessTokenManager.getIfAvailable())
);
}
}
HttpSyncDataConfiguration是Http长轮询数据同步的配置类,负责创建HttpSyncDataService(负责http数据同步的具体实现)、RestTemplate和AccessTokenManager (负责与adminhttp调用时access token的处理)。它的注解如下:
@Configuration:表示这是一个配置类;@ConditionalOnClass(HttpSyncDataService.class):条件注解,表示要有HttpSyncDataService这个类;@ConditionalOnProperty(prefix = "shenyu.sync.http", name = "url"):条件注解,要有shenyu.sync.http.url这个属性配置。@EnableConfigurationProperties(value = HttpConfig.class):表示让HttpConfig上的注解@ConfigurationProperties(prefix = "shenyu.sync.http")生效,将HttpConfig这个配置类注入Ioc容器中。
2.2 属性初始化
- HttpSyncDataService
在HttpSyncDataService的构造函数中,完成属性初始化。
public class HttpSyncDataService implements SyncDataService {
// 省略了属性字段......
public HttpSyncDataService(final HttpConfig httpConfig,
final PluginDataSubscriber pluginDataSubscriber,
final RestTemplate restTemplate,
final List<MetaDataSubscriber> metaDataSubscribers,
final List<AuthDataSubscriber> authDataSubscribers,
final AccessTokenManager accessTokenManager) {
// 1.设置accessTokenManager
this.accessTokenManager = accessTokenManager;
// 2.创建数据处理器
this.factory = new DataRefreshFactory(pluginDataSubscriber, metaDataSubscribers, authDataSubscribers);
// 3.shenyu-admin的url, 多个用逗号(,)分割
this.serverList = Lists.newArrayList(Splitter.on(",").split(httpConfig.getUrl()));
// 4.只用于http长轮询的restTemplate
this.restTemplate = restTemplate;
// 5.开始执行长轮询任务
this.start();
}
//......
}
上面代码中省略了其他函数和相 关字段,在构造函数中完成属性的初始化,主要是:
-
设置
accessTokenManager,定时向admin请求更新accessToken的值。然后每次向admin发起请求时都必须将header的X-Access-Token属性设置成accessToken对应的值; -
创建数据处理器,用于后续缓存各种类型的数据(插件、选择器、规则、元数据和认证数据);
-
获取
admin属性配置,主要是获取admin的url,admin有可能是集群,多个用逗号(,)分割; -
设置
RestTemplate,用于向admin发起请求; -
开始执行长轮询任务。
2.3 开始长轮询
- HttpSyncDataService#start()
在start()方法中,干了两件事情,一个是获取全量数据,即请求admin端获取所有需要同步的数据,然后将获取到的数据缓存到网关内存中。另一个是开启多线程执行长轮询任务。
public class HttpSyncDataService implements SyncDataService {
// ......
private void start() {
// // 只初始化一次,通过原子类实现。
if (RUNNING.compareAndSet(false, true)) {
// 初次启动,获取全量数据
this.fetchGroupConfig(ConfigGroupEnum.values());
// 一个后台服务,一个线程
int threadSize = serverList.size();
// 自定义线程池
this.executor = new ThreadPoolExecutor(threadSize, threadSize, 60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(),
ShenyuThreadFactory.create("http-long-polling", true));
// 开始长轮询,一个admin服务,创建一个线程用于数据同步
this.serverList.forEach(server -> this.executor.execute(new HttpLongPollingTask(server)));
} else {
LOG.info("shenyu http long polling was started, executor=[{}]", executor);
}
}
// ......
}
2.3.1 获取全量数据
- HttpSyncDataService#fetchGroupConfig()
ShenYu将所有需要同步的数据进行了分组,一共有5种数据类型,分别是插件、选择器、规则、元数据和认证数据。
public enum ConfigGroupEnum {
APP_AUTH, // 认证数据
PLUGIN, //插件
RULE, // 规则
SELECTOR, // 选择器
META_DATA; // 元数据
}
admin有可能是集群,这里通过循环的方式向每个admin发起请求,有一个执行成功了,那么向admin获取全量数据并缓存到网关的操作就执行成功。如果出现了异常,就向下一个admin发起请求。
public class HttpSyncDataService implements SyncDataService {
// ......
private void fetchGroupConfig(final ConfigGroupEnum... groups) throws ShenyuException {
// admin有可能是集群,这里通过循环的方式向每个admin发起请求
for (int index = 0; index < this.serverList.size(); index++) {
String server = serverList.get(index);
try {
// 真正去执行
this.doFetchGroupConfig(server, groups);
// 有一个成功,就成功了,可以退出循环
break;
} catch (ShenyuException e) {
// 出现异常,尝试执行下一个
// 最后一个也执行失败了,抛出异常
if (index >= serverList.size() - 1) {
throw e;
}
LOG.warn("fetch config fail, try another one: {}", serverList.get(index + 1));
}
}
}
// ......
}
- HttpSyncDataService#doFetchGroupConfig()
在此方法中,首先拼装请求参数,然后通过httpClient发起请求,到admin中获取数据,最后将获取到的数据更新到网关内存中。
public class HttpSyncDataService implements SyncDataService {
private void doFetchGroupConfig(final String server, final ConfigGroupEnum... groups) {
// 1. 拼请求参数,所有分组枚举类型
StringBuilder params = new StringBuilder();
for (ConfigGroupEnum groupKey : groups) {
params.append("groupKeys").append("=").append(groupKey.name()).append("&");
}
// admin端提供的接口 /configs/fetch
String url = server + Constants.SHENYU_ADMIN_PATH_CONFIGS_FETCH + "?" + StringUtils.removeEnd(params.toString(), "&");
LOG.info("request configs: [{}]", url);
String json;
try {
HttpHeaders headers = new HttpHeaders();
// 设置accessToken
headers.set(Constants.X_ACCESS_TOKEN, this.accessTokenManager.getAccessToken());
HttpEntity<String> httpEntity = new HttpEntity<>(headers);
// 2. 发起请求,获取变更数据
json = this.restTemplate.exchange(url, HttpMethod.GET, httpEntity, String.class).getBody();
} catch (RestClientException e) {
String message = String.format("fetch config fail from server[%s], %s", url, e.getMessage());
LOG.warn(message);
throw new ShenyuException(message, e);
}
// 3. 更新网关内存中数据
boolean updated = this.updateCacheWithJson(json);
if (updated) {
LOG.debug("get latest configs: [{}]", json);
return;
}
// 更新成功,此方法就执行完成了
LOG.info("The config of the server[{}] has not been updated or is out of date. Wait for 30s to listen for changes again.", server);
// 服务端没有数据更新,就等30s
ThreadUtils.sleep(TimeUnit.SECONDS, 30);
}
}
从代码中,可以看到 admin端提供的获取全量数据接口是 /configs/fetch,这里先不进一步深入,放在后文再分析。
获取到admin返回结果数据,并成功更新,那么此方法就执行结束了。如果没有更新成功,那么有可能是服务端没有数据更新,就等待30s。
这里需要提前说明一下,网关在 判断是否更新成功时,有比对数据的操作,马上就会提到。
- HttpSyncDataService#updateCacheWithJson()
更新网关内存中的数据。使用GSON进行反序列化,从属性data中拿真正的数据,然后交给DataRefreshFactory去做更新。
private boolean updateCacheWithJson(final String json) {
// 使用GSON进行反序列化
JsonObject jsonObject = GSON.fromJson(json, JsonObject.class);
// if the config cache will be updated?
return factory.executor(jsonObject.getAsJsonObject("data"));
}
- DataRefreshFactory#executor()
根据不同数据类型去更新数据,返回更新结果。具体更新逻辑交给了dataRefresh.refresh()方法。在更新结果中,有一种数据类型进行了更新,就表示此次操作发生了更新。
public boolean executor(final JsonObject data) {
//并行更新数据
List<Boolean> result = ENUM_MAP.values().parallelStream()
.map(dataRefresh -> dataRefresh.refresh(data))
.collect(Collectors.toList());
//有一个更新就表示此次发生了更新操作
return result.stream().anyMatch(Boolean.TRUE::equals);
}
- AbstractDataRefresh#refresh()
数据更新逻辑采用的是模板方法设计模式,通用操作在抽象方法中完成,不同的实现逻辑由子类完成。5种数据类型具体的更新逻辑有些差异,但是也存在通用的更新逻辑,类图关系如下:

在通用的refresh()方法中,负责数据类型转换,判断是否需要更新,和实际的数据刷新操作。
public abstract class AbstractDataRefresh<T> implements DataRefresh {
// ......
@Override
public Boolean refresh(final JsonObject data) {
// 数据类型转换
JsonObject jsonObject = convert(data);
if (Objects.isNull(jsonObject)) {
return false;
}
boolean updated = false;
// 得到数据类型
ConfigData<T> result = fromJson(jsonObject);
// 是否需要更新
if (this.updateCacheIfNeed(result)) {
updated = true;
// 真正的更新逻辑,数据刷新操作
refresh(result.getData());
}
return updated;
}
// ......
}
- AbstractDataRefresh#updateCacheIfNeed()
数据转换的过程,就是根据不同的数据类型进行转换,我们就不再进一步追踪了,看看数据是否需要更新的逻辑。方法名是updateCacheIfNeed(),通过方法重载实现。
public abstract class AbstractDataRefresh<T> implements DataRefresh {
// ......
// result是数据
protected abstract boolean updateCacheIfNeed(ConfigData<T> result);
// newVal是获取到的最新的值
// groupEnum 是哪种数据类型
protected boolean updateCacheIfNeed(final ConfigData<T> newVal, final ConfigGroupEnum groupEnum) {
// 如果是第一次,那么直接放到cache中,返回 true,表示此次进行了更新
if (GROUP_CACHE.putIfAbsent(groupEnum, newVal) == null) {
return true;
}
ResultHolder holder = new ResultHolder(false);
GROUP_CACHE.merge(groupEnum, newVal, (oldVal, value) -> {
// md5 值相同,不需要更新
if (StringUtils.equals(oldVal.getMd5(), newVal.getMd5())) {
LOG.info("Get the same config, the [{}] config cache will not be updated, md5:{}", groupEnum, oldVal.getMd5());
return oldVal;
}
// 当前缓存的数据修 改时间大于 新来的数据,不需要更新
// must compare the last update time
if (oldVal.getLastModifyTime() >= newVal.getLastModifyTime()) {
LOG.info("Last update time earlier than the current configuration, the [{}] config cache will not be updated", groupEnum);
return oldVal;
}
LOG.info("update {} config: {}", groupEnum, newVal);
holder.result = true;
return newVal;
});
return holder.result;
}
// ......
}
从上面的源码中可以看到,有两种情况不需要更新:
- 两个的数据的
md5值相同,不需要更新; - 当前缓存的数据修改时间大于 新来的数据,不需要更新。
其他情况需要更新数据。
分析到这里,就将start() 方法中初次启动,获取全量数据的逻辑分析完了,接下来是长轮询的操作。为了方 便,我将start()方法再粘贴一次:
public class HttpSyncDataService implements SyncDataService {
// ......
private void start() {
// // 只初始化一次,通过原子类实现。
if (RUNNING.compareAndSet(false, true)) {
// 初次启动,获取全量数据
this.fetchGroupConfig(ConfigGroupEnum.values());
// 一个后台服务,一个线程
int threadSize = serverList.size();
// 自定义线程池
this.executor = new ThreadPoolExecutor(threadSize, threadSize, 60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(),
ShenyuThreadFactory.create("http-long-polling", true));
// 开始长轮询,一个admin服务,创建一个线程用于数据同步
this.serverList.forEach(server -> this.executor.execute(new HttpLongPollingTask(server)));
} else {
LOG.info("shenyu http long polling was started, executor=[{}]", executor);
}
}
// ......
}
2.3.2 执行长轮询任务
- HttpLongPollingTask#run()
长轮询任务是HttpLongPollingTask,它实现了Runnable接口,任务逻辑在run()方法中。通过while()循环实现不断执行任务,即长轮询。在每一次的轮询中有三次重试逻辑,一次轮询任务失败了,等 5s 再继续,3 次都失败了,等5 分钟再试。
开始长轮询,一个admin服务,创建一个线程用于数据同步。
class HttpLongPollingTask implements Runnable {
private final String server;
HttpLongPollingTask(final String server) {
this.server = server;
}
@Override
public void run() {
// 一直轮询
while (RUNNING.get()) {
// 默认重试 3 次
int retryTimes = 3;
for (int time = 1; time <= retryTimes; time++) {
try {
doLongPolling(server);
} catch (Exception e) {
if (time < retryTimes) {
LOG.warn("Long polling failed, tried {} times, {} times left, will be suspended for a while! {}",
time, retryTimes - time, e.getMessage());
// 长轮询失败了,等 5s 再继续
ThreadUtils.sleep(TimeUnit.SECONDS, 5);
continue;
}
LOG.error("Long polling failed, try again after 5 minutes!", e);
// 3 次都失败了,等 5 分钟再试
ThreadUtils.sleep(TimeUnit.MINUTES, 5);
}
}
}
LOG.warn("Stop http long polling.");
}
}
- HttpSyncDataService#doLongPolling()
执行长轮询任务的核心逻辑:
- 根据数据类型组装请求参数:
md5和lastModifyTime; - 组装请求头和请求体;
- 向
admin发起请求,判断组数据是否发生变更; - 根据发生变更的组,再去获取数据。
public class HttpSyncDataService implements SyncDataService {
private void doLongPolling(final String server) {
// 组装请求参数:md5 和 lastModifyTime
MultiValueMap<String, String> params = new LinkedMultiValueMap<>(8);
for (ConfigGroupEnum group : ConfigGroupEnum.values()) {
ConfigData<?> cacheConfig = factory.cacheConfigData(group);
if (cacheConfig != null) {
String value = String.join(",", cacheConfig.getMd5(), String.valueOf(cacheConfig.getLastModifyTime()));
params.put(group.name(), Lists.newArrayList(value));
}
}
// 组装请求头和请求体
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_FORM_URLENCODED);
// 设置accessToken
headers.set(Constants.X_ACCESS_TOKEN, this.accessTokenManager.getAccessToken());
HttpEntity<MultiValueMap<String, String>> httpEntity = new HttpEntity<>(params, headers);
String listenerUrl = server + Constants.SHENYU_ADMIN_PATH_CONFIGS_LISTENER;
JsonArray groupJson;
//向admin发起请求,判断组数据是否发生变更
//这里只是判断了某个组是否发生变更
try {
String json = this.restTemplate.postForEntity(listenerUrl, httpEntity, String.class).getBody();
LOG.info("listener result: [{}]", json);
JsonObject responseFromServer = GsonUtils.getGson().fromJson(json, JsonObject.class);
groupJson = responseFromServer.getAsJsonArray("data");
} catch (RestClientException e) {
String message = String.format("listener configs fail, server:[%s], %s", server, e.getMessage());
throw new ShenyuException(message, e);
}
// 根据发生变更的组,再去获取数据
/**
* 官网对此处的解释:
* 网关收到响应信息之后,只知道是哪个 Group 发生了配置变更,还需要再次请求该 Group 的配置数据。
* 这里可能会存在一个疑问:为什么不是直接将变更的数据写出?
* 我们在开发的时候,也深入讨论过该问题,因为 http 长轮询机制只能保证准实时,如果在网关层处理不及时,
* 或者管理员频繁更新配置,很有可能便错过了某个配置变更的推送,安全起见,我们只告知某个 Group 信息发生了变更。
*
* 个人理解:
* 如果将变更数据直接写出,当管理员频繁更新配置时,第一次更新了,将client移除阻塞队列,返回响应信息给网关。
* 如果这个时候进行了第二次更新,那么当前的client是不在阻塞队列中,所以这一次的变更就会错过。
* 网关层处理不及时,也是同理。
* 这是一个长轮询,一个网关一个同步线程,可能存在耗时的过程。
* 如果admin有数据变更,当前网关client是没有在阻塞队列中,就不到数据。
*/
if (Objects.nonNull(groupJson) && groupJson.size() > 0) {
// fetch group configuration async.
ConfigGroupEnum[] changedGroups = GsonUtils.getGson().fromJson(groupJson, ConfigGroupEnum[].class);
LOG.info("Group config changed: {}", Arrays.toString(changedGroups));
this.doFetchGroupConfig(server, changedGroups);
}
}
}
这里需要特别解释一点的是:在长轮询任务中,为什么不直接拿到变更的数据?而是先判断哪个分组数据发生了变更,然后再次请求admin,获取变更数据?
官网对此处的解释是:
网关收到响应信息之后,只知道是哪个 Group 发生了配置变更,还需要再次请求该 Group 的配置数据。 这里可能会存在一个疑问:为什么不是直接将变更的数据写出? 我们在开发的时候,也深入讨论过该问题,因为
http长轮询机制只能保证准实时,如果在网关层处理不及时, 或者管理员频繁更新配置,很有可能便错过了某个配置变更的推送,安全起见,我们只告知某个 Group 信息发生了变更。
个人理解是:
如果将变更数据直接写出,管理员频繁更新配置时,第一次更新了,将
client