lucene-candy系列:数据同步

当前,lucene-candy实现了三种的数据同步方案

  • cn.juque.lucenecandy.core.datasync.DefaultDataSyncServiceImpl :默认的数据同步方式,仅是把数据保存到指定的索引;
  • cn.juque.lucenecandy.core.datasync.TccDataSyncServiceImpl :Try-Commit-Cancel 的方案,基于http实现数据在多实例之间的实时同步;
  • cn.juque.lucenecandy.core.datasync.MsgDataSyncServiceImpl :基于http实现多实例之间的数据批次同步,这种方案是准实时同步;

实现方案

  1. DefaultDataSyncServiceImpl
    仅是基于LuceneHelper实现的增删改查
 public void commit(IndexUpdateParamBO param) {
        switch (param.getSyncType()) {
            case ADD:
                List<Document> docList = this.documentHelper.toDocumentList(param.getContent());
                this.luceneHelper.addDocuments(param.getIndexName(), docList);
                break;
            case UPDATE:
                param.getContent().forEach(f -> {
                    BooleanQuery.Builder builder = new BooleanQuery.Builder();
                    builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.MUST);
                    this.luceneHelper.updateByQuery(param.getIndexName(), builder.build(), this.documentHelper.toDocument(f));
                });
                break;
            case DEL:
                BooleanQuery.Builder builder = new BooleanQuery.Builder();
                param.getIdList().forEach(f -> builder.add(new TermQuery(new Term(StrConstant.D_ID, f)), BooleanClause.Occur.SHOULD));
                this.luceneHelper.deleteDocuments(param.getIndexName(), builder.build());
                break;
            default:
                break;
        }
    }
  1. TccDataSyncServiceImpl
    数据同步三步走:
  • 预提交(Try):
    cn.juque.lucenecandy.core.datasync.tcc.CandyTryCommitServiceImpl 实现数据在实例间的预提交;
@Override
    public void add(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        // 保存
        IndexBO indexBO = this.documentHelper.index(entityList.get(0).getClass());
        this.luceneHelper.addDocuments(indexBO.getIndexName(), this.documentHelper.toDocumentList(entityList));
    }

    /**
     * 更新文档
     *
     * @param entityList 列表
     */
    @Override
    public void update(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        entityList.forEach(f->{
            Class<? extends BaseEntity> tClass = ClassUtil.getClass(f);
            IndexBO indexBO = this.documentHelper.index(tClass);
            // 上个版本变更为不可见
            BooleanQuery.Builder builder = new BooleanQuery.Builder();
            builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.MUST);
            int version = f.getVersion() - 1;
            builder.add(IntPoint.newExactQuery(StrConstant.D_VERSION, version), BooleanClause.Occur.MUST);
            List<Document> source = this.luceneHelper.search(indexBO.getIndexName(), builder, Sort.INDEXORDER, null);
            this.indexHelper.switchDoc(this.documentHelper.toEntityList(source, tClass), YesOrNoEnum.NO);
            // 保存新版本
            List<Document> docList = new ArrayList<>(0);
            docList.add(this.documentHelper.toDocument(f));
            this.luceneHelper.addDocuments(indexBO.getIndexName(), docList);
        });
    }

    /**
     * 删除文档
     *
     * @param className 实体类名
     * @param idList    主键列表
     */
    @Override
    public void del(String className, List<String> idList) {
        if (CollUtil.isEmpty(idList)) {
            return;
        }
        Class<? extends BaseEntity> tClass = ClassUtil.loadClass(className);
        // 根据id检索出可见文档
        BooleanQuery.Builder builder = new BooleanQuery.Builder();
        idList.forEach(f -> builder.add(new TermQuery(new Term(StrConstant.D_ID, f)), BooleanClause.Occur.SHOULD));
        IdsQueryWrapperBuilder<?> idsQueryWrapperBuilder = new IdsQueryWrapperBuilder<>(tClass, idList);
        List<? extends BaseEntity> list = this.indexHelper.searchByIds(idsQueryWrapperBuilder.build());
        // 更新为不可见
        this.indexHelper.switchDoc(list, YesOrNoEnum.NO);
    }
  • 提交(Commit):
    cn.juque.lucenecandy.core.datasync.tcc.CandyCommitServiceImpl 清除低版本数据,显现最新版本数据;
 /**
     * 添加文档
     * <li>预提交为可见文档</li>
     *
     * @param entityList 列表
     */
    @Override
    public void add(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        // 删除所有所有比当前版本低的版本
        entityList.forEach(f -> {
            IndexBO indexBO = this.documentHelper.index(f.getClass());
            int version = f.getVersion() - 1;
            BooleanQuery.Builder builder = new BooleanQuery.Builder();
            builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.MUST);
            builder.add(IntPoint.newRangeQuery(StrConstant.D_VERSION, Integer.MIN_VALUE, version), BooleanClause.Occur.MUST);
            this.luceneHelper.deleteDocuments(indexBO.getIndexName(), builder.build());
        });
    }

    /**
     * 更新文档
     *
     * @param entityList 列表
     */
    @Override
    public void update(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        IndexBO indexBO = this.documentHelper.index(entityList.get(0).getClass());
        // 移除截至上一个版本的所有版本
        BooleanQuery.Builder builder = new BooleanQuery.Builder();
        entityList.forEach(f -> {
            String id = f.getId();
            int version = f.getVersion() - 1;
            BooleanQuery.Builder delBuilder = new BooleanQuery.Builder();
            delBuilder.add(new TermQuery(new Term(StrConstant.D_ID, id)), BooleanClause.Occur.MUST);
            delBuilder.add(IntPoint.newRangeQuery(StrConstant.D_VERSION, Integer.MIN_VALUE, version), BooleanClause.Occur.MUST);
            builder.add(delBuilder.build(), BooleanClause.Occur.SHOULD);

        });
        this.luceneHelper.deleteDocuments(indexBO.getIndexName(), builder.build());
    }

    /**
     * 删除文档
     *
     * @param className 实体类名
     * @param idList    主键列表
     */
    @Override
    public void del(String className, List<String> idList) {
        if (CollUtil.isEmpty(idList)) {
            return;
        }
        // 根据id检索出文档
        BooleanQuery.Builder builder = new BooleanQuery.Builder();
        idList.forEach(f -> builder.add(new TermQuery(new Term(StrConstant.D_ID, f)), BooleanClause.Occur.SHOULD));
        Class<? extends BaseEntity> tClass = ClassUtil.loadClass(className);
        IndexBO indexBO = this.documentHelper.index(tClass);
        this.luceneHelper.deleteDocuments(indexBO.getIndexName(), builder.build());
    }

取消(Cancel):
cn.juque.lucenecandy.core.datasync.tcc.CandyCancelServiceImpl 对预提交到各实例的数据进行版本回滚;

/**
     * 添加文档
     *
     * @param entityList 列表
     */
    @Override
    public void add(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        // 根据id删除
        BooleanQuery.Builder builder = new BooleanQuery.Builder();
        entityList.forEach(f-> builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.SHOULD));
        IndexBO indexBO = this.documentHelper.index(entityList.get(0).getClass());
        this.luceneHelper.deleteDocuments(indexBO.getIndexName(), builder.build());
    }

    /**
     * 更新文档
     *
     * @param entityList 列表
     */
    @Override
    public void update(List<? extends BaseEntity> entityList) {
        if (CollUtil.isEmpty(entityList)) {
            return;
        }
        entityList.forEach(f->{
            Class<? extends BaseEntity> tClass = ClassUtil.getClass(f);
            IndexBO indexBO = this.documentHelper.index(tClass);
            // 删除当前版本
            BooleanQuery.Builder builder = new BooleanQuery.Builder();
            builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.MUST);
            builder.add(IntPoint.newExactQuery(StrConstant.D_VERSION, f.getVersion()), BooleanClause.Occur.MUST);
            this.luceneHelper.deleteDocuments(indexBO.getIndexName(), builder.build());
            // 上一个版本可见
            builder = new BooleanQuery.Builder();
            builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.MUST);
            int version = f.getVersion() - 1;
            builder.add(IntPoint.newExactQuery(StrConstant.D_VERSION, version), BooleanClause.Occur.MUST);
            List<Document> source = this.luceneHelper.search(indexBO.getIndexName(), builder, Sort.INDEXORDER, null);
            this.indexHelper.switchDoc(this.documentHelper.toEntityList(source, tClass), YesOrNoEnum.YES);
        });
    }

    /**
     * 删除文档
     *
     * @param className 实体类名
     * @param idList    主键列表
     */
    @Override
    public void del(String className, List<String> idList) {
        if (CollUtil.isEmpty(idList)) {
            return;
        }
        // 根据id检索出所有版本
        Class<? extends BaseEntity> tClass = ClassUtil.loadClass(className);
        IndexBO indexBO = this.documentHelper.index(tClass);
        BooleanQuery.Builder builder = new BooleanQuery.Builder();
        idList.forEach(f -> builder.add(new TermQuery(new Term(StrConstant.D_ID, f)), BooleanClause.Occur.SHOULD));
        List<Document> docList = this.luceneHelper.search(indexBO.getIndexName(), builder, Sort.INDEXORDER, null);
        if (CollUtil.isEmpty(docList)) {
            return;
        }
        Map<String, List<Document>> docMap = docList.stream().collect(Collectors.groupingBy(d -> d.get(StrConstant.D_ID)));
        // 过滤出最新版本,并更新为可见
        docMap.forEach((id, list) -> {
            // 根据版本号倒序
            list = list.stream().sorted(
                    Comparator.comparing(c -> Integer.parseInt(c.get(StrConstant.D_VERSION)),
                            Comparator.reverseOrder())).collect(Collectors.toList());
            // 最新版本变更为可见
            Document d = list.get(0);
            List<Document> dList = new ArrayList<>(1);
            dList.add(d);
            this.indexHelper.switchDoc(this.documentHelper.toEntityList(dList, tClass), YesOrNoEnum.YES);
            // 其他低版本统一删除
            int minVersion = Integer.parseInt(d.get(StrConstant.D_VERSION)) - 1;
            BooleanQuery.Builder delBuilder = new BooleanQuery.Builder();
            delBuilder.add(new TermQuery(new Term(StrConstant.D_ID, id)), BooleanClause.Occur.MUST);
            delBuilder.add(IntPoint.newRangeQuery(StrConstant.D_VERSION, Integer.MIN_VALUE, minVersion), BooleanClause.Occur.MUST);
            this.luceneHelper.deleteDocuments(indexBO.getIndexName(), delBuilder.build());
        });
    }
  • TccDataSyncServiceImpl
    实现父类IDataSyncService,把Tcc集成到IndexHelper。
@Slf4j
@Service("tccDataSyncService")
@DataSync(value = "tcc")
public class TccDataSyncServiceImpl implements IDataSyncService {

    ...

    @Resource(name = "candyTryCommitService")
    private ITccService candyTryCommitService;

    @Resource(name = "candyCommitService")
    private ITccService candyCommitService;

    @Resource(name = "candyCancelService")
    private ITccService candyCancelService;

    /**
     * 更新
     *
     * @param param 数据
     */
    @Override
    public void commit(@NotNull IndexUpdateParamBO param) {
        // 预提交
        boolean flag = this.tryOperate(param);
        // 实际没有对其他节点执行预提交 或 提交成功
        if (flag) {
            this.commitOperate(param);
        } else {
            // 取消
            this.cancel(param);
            throw new AppException(MessageEnum.SYSTEM_ERROR);
        }
    }

    /**
     * 删除
     *
     * @param param 数据
     */
    @Override
    public void cancel(@NotNull IndexUpdateParamBO param) {
        switch (param.getSyncType()) {
            case ADD:
                this.candyCancelService.add(param.getContent());
                break;
            case UPDATE:
                this.candyCancelService.update(param.getContent());
                break;
            case DEL:
                this.candyCancelService.del(param.getClassName(), param.getIdList());
                break;
            default:
                break;
        }
        // cancel
        this.batchPost(param, SyncTypeEnum.CANCEL);
    }
}
  1. MsgDataSyncServiceImpl
    数据的批次准时同步逻辑:
  • 收集数据操作,IndexUpdateParamBO 封装了一次数据操作,并存放到本地队列
public void commit(IndexUpdateParamBO param) {
        param.setTimestamp(System.currentTimeMillis());
        this.dataSyncService.commit(param);
        if(!this.commitListenerRender.before(param)) {
            return;
        }
        this.dataSyncWaitReadMsgCache.put(null, param);
        this.commitListenerRender.after(param);
    }
  • 队列数据写入到本地存储,即:_candy_data_sync_message
private void writeToDisk() {
        if (!TAKE_END.compareAndSet(true, false)) {
            return;
        }
        String className = ClassUtil.getClassName(BaseMessageEntity.class, false);
        IndexBO indexBO = this.indexInfoCache.get(className, false);
        List<IndexUpdateParamBO> list = new ArrayList<>(100);
        IndexUpdateParamBO paramBO;
        try {
            boolean flag = true;
            while (flag) {
                paramBO = WAIT_READ.poll();
                if (Objects.isNull(paramBO)) {
                    flag = false;
                    continue;
                }
                list.add(paramBO);
                if (list.size() >= 99) {
                    this.saveLog(indexBO, list);
                    list.clear();
                }
            }
            if (CollUtil.isNotEmpty(list)) {
                this.saveLog(indexBO, list);
            }
        } catch (Exception e) {
            log.error("msg write to disk error", e);
            Thread.currentThread().interrupt();
        } finally {
            TAKE_END.set(true);
        }
    }

至此,完成数据操作的本地持久化。接下来,由DataSyncMsgReadThreadPool完成数据操作到实例的同步。该异步线程会进行批量同步,并在同步结束后移除本地持久化。

private void execute() {
        if (!IS_END.compareAndSet(true, false)) {
            return;
        }
        String className = ClassUtil.getClassName(BaseMessageEntity.class, false);
        IndexBO indexBO = this.indexInfoCache.get(className, false);
        String indexName = indexBO.getIndexName();
        try {
            List<Document> docList = this.luceneHelper.searchByPage(indexName, BOOL_BUILDER, SORT, PAGE_INFO, null);
            while (CollUtil.isNotEmpty(docList)) {
                List<BaseMessageEntity> entityList = this.documentHelper.toEntityList(docList, BaseMessageEntity.class);
                this.batchPost(entityList);
                // 删除已同步的文档
                BooleanQuery.Builder builder = new BooleanQuery.Builder();
                entityList.forEach(
                        f -> builder.add(new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.SHOULD));
                this.luceneHelper.deleteDocuments(indexName, builder.build());
                docList = this.luceneHelper.searchByPage(indexName, BOOL_BUILDER, SORT, PAGE_INFO, null);
            }
        } catch (Exception e) {
            log.error("文档同步异常", e);
        } finally {
            IS_END.set(true);
        }
    }

实例接收到数据操作报文后,同步到本地索引:_candy_data_sync_message。最后由DataSyncMsgWriteThreadPool 解析并同步到对应的索引。

private void execute() {
        if (!IS_END.compareAndSet(true, false)) {
            return;
        }
        String className = ClassUtil.getClassName(BaseMessageEntity.class, false);
        IndexBO indexBO = this.indexInfoCache.get(className, false);
        String indexName = indexBO.getIndexName();
        try {
            List<Document> docList = this.luceneHelper.searchByPage(indexName, BOOL_BUILDER, SORT, PAGE_INFO, null);
            while (CollUtil.isNotEmpty(docList)) {
                List<BaseMessageEntity> entityList = this.documentHelper.toEntityList(docList, BaseMessageEntity.class);
                this.write(entityList);
                // 删除已同步的文档
                BooleanQuery.Builder builder = new BooleanQuery.Builder();
                entityList.forEach(f ->
                        builder.add(
                                new TermQuery(new Term(StrConstant.D_ID, f.getId())), BooleanClause.Occur.SHOULD));
                this.luceneHelper.deleteDocuments(indexName, builder.build());
                docList = this.luceneHelper.searchByPage(indexName, BOOL_BUILDER, SORT, PAGE_INFO, null);
            }
        } catch (Exception e) {
            log.error("文档同步异常", e);
        } finally {
            IS_END.set(true);
        }
    }

至此,消息同步数据完成。

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 218,451评论 6 506
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 93,172评论 3 394
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 164,782评论 0 354
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,709评论 1 294
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,733评论 6 392
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,578评论 1 305
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,320评论 3 418
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 39,241评论 0 276
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,686评论 1 314
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,878评论 3 336
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,992评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,715评论 5 346
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,336评论 3 330
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,912评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 33,040评论 1 270
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 48,173评论 3 370
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,947评论 2 355

推荐阅读更多精彩内容