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

做官网网站哪家公司好新手开网店从哪里找货源

做官网网站哪家公司好,新手开网店从哪里找货源,app开发与网站建设难度,网站首页被k怎么恢复背景 flink在实现本地内存和db同步配置表信息时,想要做到类似于增量(保证实时性) 全量(保证和DB数据一致)的效果,那么我们如何通过flink的广播状态外部定时器定时全量同步的方式来实现呢? 实现增量全量的效果 package wikiedits.schedule…

背景

flink在实现本地内存和db同步配置表信息时,想要做到类似于增量(保证实时性) + 全量(保证和DB数据一致)的效果,那么我们如何通过flink的广播状态+外部定时器定时全量同步的方式来实现呢?

实现增量+全量的效果

package wikiedits.schedule;import java.util.List;
import java.util.Map;import org.apache.commons.lang3.StringUtils;
import org.apache.flink.api.common.state.BroadcastState;
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
import org.apache.flink.api.java.typeutils.ListTypeInfo;
import org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction;
import org.apache.flink.util.Collector;//处理函数
public class BroadcastStatePlusSchedulerFunction extends KeyedBroadcastProcessFunction<String, String, String, String> {// 键值分区状态private final MapStateDescriptor<String, List<String>> mapStateDesc =new MapStateDescriptor<>("items", BasicTypeInfo.STRING_TYPE_INFO, new ListTypeInfo<>(String.class));// 广播状态private final MapStateDescriptor<String, String> ruleStateDescriptor = new MapStateDescriptor<>("RulesBroadcastState", BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO);@Overridepublic void processBroadcastElement(String value, Context ctx, Collector<String> out) throws Exception {// 1.增量消息更新广播状态BroadcastState<String, String> broadcastState = ctx.getBroadcastState(ruleStateDescriptor);broadcastState.put(value, value);// 2.全量更新,判断广播状态和DB配置表在本地缓存的配置项是否一致,比如如果广播状态记录少了,使用本地缓存中的记录来更新下广播状态for (Map.Entry<String, String> entry : StaticLoadUtil.getConfigCache().asMap().entrySet()) {String broadcastValue = broadcastState.get(entry.getKey());if(!StringUtils.equals(entry.getValue(), broadcastValue)){//如果不相等,那么以DB缓存中的为准}}// 3.自此,广播状态和DB配置表的状态几乎一致,不过由于他们的比较只发生于收到广播元素,所以我们可以在凌晨的时候故意从db中找出几条记录发送kafka消息到这个广播状态来进行触发比较,当然这里也可以当收到某个元素时覆盖掉flink的广播状态}@Overridepublic void processElement(String value, ReadOnlyContext ctx, Collector<String> out) throws Exception {// 键值分区状态final MapState<String, List<String>> state = getRuntimeContext().getMapState(mapStateDesc);// 广播状态for (Map.Entry<String, String> entry : ctx.getBroadcastState(ruleStateDescriptor).immutableEntries()) {}}}// 外部定时器实现
package wikiedits.schedule;import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;import com.google.common.cache.Cache;
import com.google.common.cache.CacheBuilder;/*** 静态类定时加载DB配置表到本地内存中*/
public class StaticLoadUtil {// 定时任务执行器private static transient ScheduledExecutorService scheduledExecutorService;public static final Cache<String, String> configCache =CacheBuilder.newBuilder().initialCapacity(50).maximumSize(500).build();// 通过定时执行器定时同步本地缓存和DB配置表static {scheduledExecutorService = Executors.newScheduledThreadPool(10);scheduledExecutorService.scheduleWithFixedDelay(() -> {// 2.1 定时任务更新本地内存配置项// List<ConfigEntity> configList = DBManager.SELECTSQL.getConfigs();// for(ConfigEntity entity : configList){configCache.put("key", "value");// }// 2.2 更新本地变量threshold的值// threshold = DBManager.SELECTSQL.getConfig("threshold");}, 0, 100, TimeUnit.SECONDS);}/*** 获取本地缓存*/public static Cache<String, String> getConfigCache() {return configCache;}}

总结:

1.在处理广播元素的时候,除了更新广播状态之外,还要对比下广播状态和DB配置表在flink的本地缓存的数据,如果不一致,需要打印告警日志或者采取更新等措施

2.由于全量广播状态和DB配置表在flink的本地缓存的数据对比是在接收到某个广播元素的时候才进行,所以我们可以多余多发送一些相同的广播元素来触发对比

3.通过这种方式,广播状态就可以实现增量(实时性) + 全量(准确性) 的结果

http://www.sczhlp.com/news/140716/

相关文章:

  • 珠海网站建设q479185700棒西安软件优化网站建设
  • 正规的培训行业网站开发青海公路建设信息服务网站
  • 企业微网站建设方案网站建设业务培训资料
  • 黄页网站营销软件开发外包商业模式
  • 网站前端设计与实现网络工程是冷门专业吗
  • 空白金兰契的多维解构与实践路径:从价值表征困境到人机共生伦理
  • 溧阳做网站哪家好那个网站的详情做的好
  • 太平阳电脑网网站模板企石网站建设
  • 长春省妇幼网站做四维asp网站301
  • 潮品服饰网站建设规划书网站的建设合同是否交印花税
  • 能自己在家做网站吗北京哪家网站建设公司好
  • wordpress 站群插件同城招聘网站自助建站
  • 医院网站建设联系方式如何做网站后台
  • 做土特产网站什么名字最好环保部网站建设项目重大变动
  • 沈阳成创网站建设公司哪个网站可以做计算机二级的题
  • 如何本地搭建自己的网站后台网站模板下载
  • 2025中国制造企业500强榜单发布
  • 珠海高端网站设计网站建设和推广话术6
  • 宿迁房产网官方网站温州网站建设制作公司
  • 有关网站建设的文章设计网站导航大全
  • 手机网站编程大良做网站的公司
  • 读 WPF 源代码 了解获取 GlyphTypeface 的 CharacterToGlyphMap 的数量耗时原因
  • 张江,首个万亿市值巨头诞生!
  • 怎么创立网站page转wordpress
  • 网页设计需求模板seo技巧是什么
  • 广州网站建设网络西安网站开发联系方式
  • 模版 网站需要多少钱wordpress网页设计步骤
  • 网页内嵌网站域名和网站
  • 友情链接交易网站源码网站开发需要如何压缩代码
  • 中国工商银行官方网站登录建设职业技术学院网站