This is an automated email from the ASF dual-hosted git repository.

aiceflower pushed a commit to branch release-0.9.4
in repository https://gitbox.apache.org/repos/asf/linkis.git

commit 6bf6551f51fddf6a737c88401966634d53c3a1d9
Author: chaogefeng <[email protected]>
AuthorDate: Wed Jun 3 18:03:23 2020 +0800

    To add support for CS Listener
    close #374
---
 contextservice/cs-listener/pom.xml                 |  71 ++++++++++
 .../linkis/cs/listener/CSIDListener.java           |  19 +++
 .../linkis/cs/listener/CSKeyListener.java          |  16 +++
 .../cs/listener/ContextAsyncEventListener.java     |  25 ++++
 .../ListenerBus/ContextAsyncListenerBus.java       |  40 ++++++
 .../listener/callback/AbstractCallbackEngine.java  |  12 ++
 .../cs/listener/callback/CallbackEngine.java       |  16 +++
 .../listener/callback/ContextIDCallbackEngine.java |  12 ++
 .../callback/ContextKeyCallbackEngine.java         |  11 ++
 .../listener/callback/imp/ContextKeyValueBean.java |  57 ++++++++
 .../imp/DefaultContextIDCallbackEngine.java        | 149 ++++++++++++++++++++
 .../imp/DefaultContextKeyCallbackEngine.java       | 153 +++++++++++++++++++++
 .../cs/listener/conf/ContextListenerConf.java      |  13 ++
 .../linkis/cs/listener/event/ContextIDEvent.java   |  13 ++
 .../linkis/cs/listener/event/ContextKeyEvent.java  |  10 ++
 .../cs/listener/event/enumeration/OperateType.java |  13 ++
 .../listener/event/impl/DefaultContextIDEvent.java |  34 +++++
 .../event/impl/DefaultContextKeyEvent.java         |  56 ++++++++
 .../cs/listener/manager/ListenerManager.java       |  15 ++
 .../manager/imp/DefaultContextListenerManager.java |  46 +++++++
 .../linkis/cs/listener/test/TestContextID.java     |  22 +++
 .../linkis/cs/listener/test/TestContextKey.java    |  63 +++++++++
 .../cs/listener/test/TestContextKeyValue.java      |  36 +++++
 .../linkis/cs/listener/test/TestContextValue.java  |  37 +++++
 .../cs/listener/test/TestListenerManager.java      | 108 +++++++++++++++
 25 files changed, 1047 insertions(+)

diff --git a/contextservice/cs-listener/pom.xml 
b/contextservice/cs-listener/pom.xml
new file mode 100644
index 0000000000..18c3482a9b
--- /dev/null
+++ b/contextservice/cs-listener/pom.xml
@@ -0,0 +1,71 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+  ~ Copyright 2019 WeBank
+  ~ Licensed under the Apache License, Version 2.0 (the "License");
+  ~  you may not use this file except in compliance with the License.
+  ~ You may obtain a copy of the License at
+  ~ http://www.apache.org/licenses/LICENSE-2.0
+  ~ Unless required by applicable law or agreed to in writing, software
+  ~ distributed under the License is distributed on an "AS IS" BASIS,
+  ~ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+  ~ See the License for the specific language governing permissions and
+  ~ limitations under the License.
+  -->
+
+<project xmlns="http://maven.apache.org/POM/4.0.0";
+         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance";
+         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <parent>
+        <artifactId>linkis</artifactId>
+        <groupId>com.webank.wedatasphere.linkis</groupId>
+        <version>0.9.4</version>
+        <relativePath>../../pom.xml</relativePath>
+    </parent>
+    <modelVersion>4.0.0</modelVersion>
+
+    <artifactId>linkis-cs-listener</artifactId>
+
+    <dependencies>
+        <dependency>
+            <groupId>com.webank.wedatasphere.linkis</groupId>
+            <artifactId>linkis-common</artifactId>
+            <version>${linkis.version}</version>
+            <scope>provided</scope>
+        </dependency>
+        <dependency>
+            <groupId>com.webank.wedatasphere.linkis</groupId>
+            <artifactId>linkis-cs-common</artifactId>
+            <version>0.9.4</version>
+        </dependency>
+        <dependency>
+            <groupId>junit</groupId>
+            <artifactId>junit</artifactId>
+            <version>4.12</version>
+            <scope>test</scope>
+        </dependency>
+    </dependencies>
+
+    <build>
+        <plugins>
+            <plugin>
+                <groupId>org.apache.maven.plugins</groupId>
+                <artifactId>maven-deploy-plugin</artifactId>
+            </plugin>
+
+            <plugin>
+                <groupId>net.alchim31.maven</groupId>
+                <artifactId>scala-maven-plugin</artifactId>
+            </plugin>
+            <plugin>
+                <groupId>org.apache.maven.plugins</groupId>
+                <artifactId>maven-jar-plugin</artifactId>
+            </plugin>
+        </plugins>
+        <resources>
+            <resource>
+                <directory>${basedir}/src/main/resources</directory>
+            </resource>
+        </resources>
+        <finalName>${project.artifactId}-${project.version}</finalName>
+    </build>
+</project>
\ No newline at end of file
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSIDListener.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSIDListener.java
new file mode 100644
index 0000000000..458ed83a3c
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSIDListener.java
@@ -0,0 +1,19 @@
+package com.webank.wedatasphere.linkis.cs.listener;
+
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextIDEvent;
+
+/**
+ * @author peacewong
+ * @date 2020/2/15 11:24
+ */
+public interface CSIDListener extends ContextAsyncEventListener {
+
+
+
+    void onCSIDAccess(ContextIDEvent contextIDEvent);
+
+    void onCSIDADD(ContextIDEvent contextIDEvent);
+
+    void onCSIDRemoved(ContextIDEvent contextIDEvent);
+
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSKeyListener.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSKeyListener.java
new file mode 100644
index 0000000000..13707bdad0
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/CSKeyListener.java
@@ -0,0 +1,16 @@
+package com.webank.wedatasphere.linkis.cs.listener;
+
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextKeyEvent;
+
+/**
+ * @author peacewong
+ * @date 2020/2/15 11:24
+ */
+public interface CSKeyListener extends ContextAsyncEventListener {
+    @Override
+    void onEvent(Event event);
+    void onCSKeyUpdate(ContextKeyEvent contextKeyEvent);
+    void onCSKeyAccess(ContextKeyEvent contextKeyEvent);
+
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ContextAsyncEventListener.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ContextAsyncEventListener.java
new file mode 100644
index 0000000000..008e20216b
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ContextAsyncEventListener.java
@@ -0,0 +1,25 @@
+/*
+ * Copyright 2019 WeBank
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ *  you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ * http://www.apache.org/licenses/LICENSE-2.0
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package com.webank.wedatasphere.linkis.cs.listener;
+
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import com.webank.wedatasphere.linkis.common.listener.EventListener;
+
+/**
+ * @author peacewong
+ * @date 2020/2/22 21:45
+ */
+public interface ContextAsyncEventListener extends EventListener {
+    void onEvent(Event event);
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ListenerBus/ContextAsyncListenerBus.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ListenerBus/ContextAsyncListenerBus.java
new file mode 100644
index 0000000000..5b16dfca76
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/ListenerBus/ContextAsyncListenerBus.java
@@ -0,0 +1,40 @@
+package com.webank.wedatasphere.linkis.cs.listener.ListenerBus;
+
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import com.webank.wedatasphere.linkis.common.listener.ListenerEventBus;
+import com.webank.wedatasphere.linkis.cs.listener.ContextAsyncEventListener;
+import com.webank.wedatasphere.linkis.cs.listener.conf.ContextListenerConf;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/21
+ */
+public class ContextAsyncListenerBus<L extends ContextAsyncEventListener, E 
extends Event> extends ListenerEventBus<L, E> {
+
+
+    private static final String NAME = "ContextAsyncListenerBus";
+
+    public ContextAsyncListenerBus() {
+        super(ContextListenerConf.WDS_CS_LISTENER_ASYN_QUEUE_CAPACITY, NAME, 
ContextListenerConf.WDS_CS_LISTENER_ASYN_CONSUMER_THREAD_MAX, 
ContextListenerConf.WDS_CS_LISTENER_ASYN_CONSUMER_THREAD_FREE_TIME_MAX);
+    }
+
+    @Override
+    public void doPostEvent(L listener, E event) {
+        listener.onEvent(event);
+    }
+
+
+    private static ContextAsyncListenerBus contextAsyncListenerBus = null;
+
+    public static ContextAsyncListenerBus getInstance() {
+        if (contextAsyncListenerBus == null) {
+            synchronized (ContextAsyncListenerBus.class) {
+                if (contextAsyncListenerBus == null) {
+                    contextAsyncListenerBus = new ContextAsyncListenerBus();
+                    contextAsyncListenerBus.start();
+                }
+            }
+        }
+        return contextAsyncListenerBus;
+    }
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/AbstractCallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/AbstractCallbackEngine.java
new file mode 100644
index 0000000000..5a27b03caa
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/AbstractCallbackEngine.java
@@ -0,0 +1,12 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ */
+public interface AbstractCallbackEngine extends CallbackEngine {
+    //todo
+    //实现事件的存储和按需消费:存储这些变化的事件,并且按需消费
+    //事件超过一定时间还没被消费,自动移除
+    //cskey被五个client注册了listener,如果有挂掉,那么必须要一个最大消费时间的机制
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/CallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/CallbackEngine.java
new file mode 100644
index 0000000000..04b059ff90
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/CallbackEngine.java
@@ -0,0 +1,16 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback;
+
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.ContextKeyValueBean;
+
+import java.util.ArrayList;
+import java.util.Set;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ */
+public interface CallbackEngine {
+
+    ArrayList<ContextKeyValueBean> getListenerCallback(String source);
+
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextIDCallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextIDCallbackEngine.java
new file mode 100644
index 0000000000..f5532c2d56
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextIDCallbackEngine.java
@@ -0,0 +1,12 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.listener.ListenerDomain;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ *
+ */
+public interface ContextIDCallbackEngine extends CallbackEngine {
+    void registerClient(ListenerDomain listenerDomain);
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextKeyCallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextKeyCallbackEngine.java
new file mode 100644
index 0000000000..48f4c89bbf
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/ContextKeyCallbackEngine.java
@@ -0,0 +1,11 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.listener.ListenerDomain;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ */
+public interface ContextKeyCallbackEngine extends CallbackEngine {
+    void registerClient(ListenerDomain listenerDomain);
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/ContextKeyValueBean.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/ContextKeyValueBean.java
new file mode 100644
index 0000000000..d9fa5f57f5
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/ContextKeyValueBean.java
@@ -0,0 +1,57 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback.imp;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextValue;
+
+import java.util.Objects;
+
+public class ContextKeyValueBean {
+
+    private ContextKey csKey;
+    private ContextValue csValue;
+    private ContextID csID;
+
+    public ContextID getCsID() {
+        return csID;
+    }
+
+    public void setCsID(ContextID csID) {
+        this.csID = csID;
+    }
+
+    public ContextKey getCsKey() {
+        return csKey;
+    }
+
+    public void setCsKey(ContextKey csKey) {
+        this.csKey = csKey;
+    }
+
+    public ContextValue getCsValue() {
+        return csValue;
+    }
+
+    public void setCsValue(ContextValue csValue) {
+        this.csValue = csValue;
+    }
+
+
+    @Override
+    public int hashCode() {
+        return Objects.hash(csKey, csValue);
+    }
+
+    @Override
+    public boolean equals(Object o) {
+        if (this == o) {
+            return true;
+        }
+        if (o == null || getClass() != o.getClass()) {
+            return false;
+        }
+        ContextKeyValueBean csmapKey = (ContextKeyValueBean) o;
+        return Objects.equals(csKey, csmapKey.csKey) &&
+                Objects.equals(csValue, csmapKey.csValue);
+    }
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextIDCallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextIDCallbackEngine.java
new file mode 100644
index 0000000000..5008002cea
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextIDCallbackEngine.java
@@ -0,0 +1,149 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback.imp;
+
+import com.google.common.collect.HashMultimap;
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import 
com.webank.wedatasphere.linkis.cs.common.entity.listener.CommonContextIDListenerDomain;
+import com.webank.wedatasphere.linkis.cs.common.entity.listener.ListenerDomain;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+import com.webank.wedatasphere.linkis.cs.common.listener.ContextIDListener;
+import com.webank.wedatasphere.linkis.cs.listener.CSIDListener;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.ContextIDCallbackEngine;
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextIDEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextIDEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextKeyEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.manager.imp.DefaultContextListenerManager;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+
+import static 
com.webank.wedatasphere.linkis.cs.listener.event.enumeration.OperateType.ADD;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ */
+public class DefaultContextIDCallbackEngine implements CSIDListener, 
ContextIDCallbackEngine {
+
+    private static final Logger logger = 
LoggerFactory.getLogger(DefaultContextIDCallbackEngine.class);
+    private HashMultimap<String, ContextID> registerCSIDcsClients = 
HashMultimap.create();//key为clientSource的instance值
+
+    private List<ContextID> removedContextIDS = new ArrayList<>();
+
+
+    @Override
+    public ArrayList<ContextKeyValueBean> getListenerCallback(String source) {
+        Set<ContextID> ContextIDSets = registerCSIDcsClients.get(source);
+        ArrayList<ContextKeyValueBean> contextKeyValueBeans = new 
ArrayList<>();
+        for (ContextID contextID : removedContextIDS) {
+            if (ContextIDSets.contains(contextID)) {
+                ContextKeyValueBean contextKeyValueBean = new 
ContextKeyValueBean();
+                contextKeyValueBean.setCsID(contextID);
+                contextKeyValueBeans.add(contextKeyValueBean);
+            }
+        }
+        return contextKeyValueBeans;
+    }
+
+
+    @Override
+    public void registerClient(ListenerDomain listenerDomain) {
+        if (listenerDomain != null && listenerDomain instanceof 
CommonContextIDListenerDomain) {
+            CommonContextIDListenerDomain commonContextIDListenerDomain = 
(CommonContextIDListenerDomain) listenerDomain;
+            String source = commonContextIDListenerDomain.getSource();
+            ContextID contextID = commonContextIDListenerDomain.getContextID();
+            if (source != null && contextID != null) {
+                synchronized (registerCSIDcsClients) {
+                    registerCSIDcsClients.put(source, contextID);
+                }
+            }
+        }
+    }
+
+        @Override
+        public void onEvent (Event event){
+            DefaultContextIDEvent defaultContextIDEvent = null;
+            if (event != null && event instanceof DefaultContextIDEvent) {
+                defaultContextIDEvent = (DefaultContextIDEvent) event;
+            }
+            if (null == defaultContextIDEvent) {
+                logger.warn("defaultContextIDEvent event 为空");
+                return;
+            }
+            switch (defaultContextIDEvent.getOperateType()) {
+                //ADD, UPDATE, DELETE, REMOVEALL, ACCESS
+                case REMOVEALL:
+                    onCSIDRemoved(defaultContextIDEvent);
+                    break;
+                case ADD:
+                    onCSIDADD(defaultContextIDEvent);
+                    break;
+                case ACCESS:
+                    onCSIDAccess(defaultContextIDEvent);
+                    break;
+                case UPDATE:
+                    break;
+                case DELETE:
+                    break;
+                default:
+                    logger.info("检查defaultContextIDEvent event操作类型");
+            }
+
+        }
+
+        @Override
+        public void onCSIDAccess (ContextIDEvent contextIDEvent){
+
+        }
+
+        @Override
+        public void onCSIDADD (ContextIDEvent contextIDEvent){
+
+        }
+
+        @Override
+        public void onCSIDRemoved (ContextIDEvent contextIDEvent){
+
+            DefaultContextIDEvent defaultContextIDEvent = null;
+            if (contextIDEvent != null && contextIDEvent instanceof 
DefaultContextIDEvent) {
+                defaultContextIDEvent = (DefaultContextIDEvent) contextIDEvent;
+            }
+            if (null == defaultContextIDEvent) {
+                return;
+            }
+            synchronized (removedContextIDS) {
+                removedContextIDS.add(defaultContextIDEvent.getContextID());
+            }
+        }
+
+        @Override
+        public void onEventError (Event event, Throwable t){
+
+        }
+
+
+        private static DefaultContextIDCallbackEngine 
singleDefaultContextIDCallbackEngine = null;
+
+        private DefaultContextIDCallbackEngine() {
+
+        }
+
+        public static DefaultContextIDCallbackEngine getInstance () {
+            if (singleDefaultContextIDCallbackEngine == null) {
+                synchronized (DefaultContextIDCallbackEngine.class) {
+                    if (singleDefaultContextIDCallbackEngine == null) {
+                        singleDefaultContextIDCallbackEngine = new 
DefaultContextIDCallbackEngine();
+                        DefaultContextListenerManager 
instanceContextListenerManager = DefaultContextListenerManager.getInstance();
+                        
instanceContextListenerManager.getContextAsyncListenerBus().addListener(singleDefaultContextIDCallbackEngine);
+                        logger.info("add listerner 
singleDefaultContextIDCallbackEngine success");
+                    }
+                }
+            }
+            return singleDefaultContextIDCallbackEngine;
+        }
+    }
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextKeyCallbackEngine.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextKeyCallbackEngine.java
new file mode 100644
index 0000000000..154bed89b8
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/callback/imp/DefaultContextKeyCallbackEngine.java
@@ -0,0 +1,153 @@
+package com.webank.wedatasphere.linkis.cs.listener.callback.imp;
+
+import com.google.common.collect.HashMultimap;
+import com.google.common.collect.Multiset;
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import 
com.webank.wedatasphere.linkis.cs.common.entity.listener.CommonContextIDListenerDomain;
+import 
com.webank.wedatasphere.linkis.cs.common.entity.listener.CommonContextKeyListenerDomain;
+import com.webank.wedatasphere.linkis.cs.common.entity.listener.ListenerDomain;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+import com.webank.wedatasphere.linkis.cs.listener.CSKeyListener;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.ContextKeyCallbackEngine;
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextKeyEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextIDEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextKeyEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.manager.imp.DefaultContextListenerManager;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.*;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/20
+ */
+public class DefaultContextKeyCallbackEngine implements CSKeyListener, 
ContextKeyCallbackEngine {
+    private static final Logger logger = 
LoggerFactory.getLogger(DefaultContextKeyCallbackEngine.class);
+
+    private HashMultimap<String, ContextID> registerCSIDcsClients = 
HashMultimap.create();//key为clientSource的instance值
+
+    private HashMultimap<String, ContextKeyValueBean> registerCSIDcsKeyValues 
= HashMultimap.create();//key为 contextId 的ID值
+
+    //注册csClient及其监听的csKeys
+    @Override
+    public void registerClient(ListenerDomain listenerDomain) {
+        if (listenerDomain != null && listenerDomain instanceof 
CommonContextKeyListenerDomain) {
+            CommonContextKeyListenerDomain commonContextKeyListenerDomain = 
(CommonContextKeyListenerDomain) listenerDomain;
+            String source = commonContextKeyListenerDomain.getSource();
+            ContextID contextID = 
commonContextKeyListenerDomain.getContextID();
+            ContextKey contextKey = 
commonContextKeyListenerDomain.getContextKey();
+            if (source != null && contextID != null) {
+                synchronized (registerCSIDcsClients) {
+                    logger.info("要注册的csClient和contextId: " + source + ":" + 
contextID);
+                    registerCSIDcsClients.put(source, contextID);
+                }
+            }
+            //针对cskey生成一个bean,cskey对应的value值目前为空
+            if (contextKey != null) {
+                ContextKeyValueBean contextKeyValueBean = new 
ContextKeyValueBean();
+                contextKeyValueBean.setCsKey(contextKey);
+                contextKeyValueBean.setCsID(contextID);
+                synchronized (registerCSIDcsKeyValues) {
+                    logger.info("要注册的contextId: " + contextID.getContextId());
+                    registerCSIDcsKeyValues.put(contextID.getContextId(), 
contextKeyValueBean);
+                }
+            }
+        }
+    }
+
+    //通过 source 拿到 ContextID,遍历 ContextID,返回监听的 beans
+    @Override
+    public ArrayList<ContextKeyValueBean> getListenerCallback(String source) {
+        ArrayList<ContextKeyValueBean> arrayContextKeyValueBeans = new 
ArrayList<>();
+        Set<ContextID> contextIDS = registerCSIDcsClients.get(source);
+        //返回所有的 ContextKeyValueBean
+        if (contextIDS.size() > 0) {
+            for (ContextID csId : contextIDS) {
+                
arrayContextKeyValueBeans.addAll(registerCSIDcsKeyValues.get(csId.getContextId()));
+            }
+        }
+        return arrayContextKeyValueBeans;
+    }
+
+    @Override
+    public void onEvent(Event event) {
+        DefaultContextKeyEvent defaultContextKeyEvent = null;
+        if (event != null && event instanceof DefaultContextKeyEvent) {
+            defaultContextKeyEvent = (DefaultContextKeyEvent) event;
+        }
+        if (null == defaultContextKeyEvent) {
+            logger.info("defaultContextKeyEvent event 为空");
+            return;
+        }
+        logger.info("defaultContextKeyEvent 要更新事件的ID: " + 
defaultContextKeyEvent.getContextID().getContextId());
+        logger.info("defaultContextKeyEvent 要更新事件的key: " + 
defaultContextKeyEvent.getContextKeyValue().getContextKey().getKey());
+        logger.info("defaultContextKeyEvent 要更新的value" + 
defaultContextKeyEvent.getContextKeyValue().getContextValue().getValue());
+        switch (defaultContextKeyEvent.getOperateType()) {
+            case UPDATE:
+                onCSKeyUpdate(defaultContextKeyEvent);
+                break;
+            case ACCESS:
+                onCSKeyAccess(defaultContextKeyEvent);
+                break;
+            default:
+                logger.info("检查defaultContextKeyEvent event操作类型");
+        }
+    }
+
+    //更新 cskey 对应的 value 值
+    @Override
+    public void onCSKeyUpdate(ContextKeyEvent cskeyEvent) {
+        DefaultContextKeyEvent defaultContextKeyEvent = null;
+        if (cskeyEvent != null && cskeyEvent instanceof 
DefaultContextKeyEvent) {
+            defaultContextKeyEvent = (DefaultContextKeyEvent) cskeyEvent;
+        }
+        if (null == defaultContextKeyEvent) {
+            return;
+        }
+
+        synchronized (registerCSIDcsKeyValues) {
+            //遍历所有csid,如果csid跟事件中的相同,则取出该csid所有的bean,更新所有bean中的csvalue.
+            Set<ContextKeyValueBean> contextKeyValueBeans = 
registerCSIDcsKeyValues.get(defaultContextKeyEvent.getContextID().getContextId());
+            for (ContextKeyValueBean contextKeyValueBean : 
contextKeyValueBeans) {
+                if 
(contextKeyValueBean.getCsKey().getKey().equals(defaultContextKeyEvent.getContextKeyValue().getContextKey().getKey()))
 {
+                    
contextKeyValueBean.setCsValue(defaultContextKeyEvent.getContextKeyValue().getContextValue());
+                }
+            }
+        }
+    }
+
+
+    //todo
+    @Override
+    public void onCSKeyAccess(ContextKeyEvent cskeyEvent) {
+
+    }
+
+    //todo
+    @Override
+    public void onEventError(Event event, Throwable t) {
+
+    }
+
+    private static DefaultContextKeyCallbackEngine 
singleDefaultContextKeyCallbackEngine = null;
+
+    private DefaultContextKeyCallbackEngine() {
+    }
+
+    public static DefaultContextKeyCallbackEngine getInstance() {
+        if (singleDefaultContextKeyCallbackEngine == null) {
+            synchronized (DefaultContextKeyCallbackEngine.class) {
+                if (singleDefaultContextKeyCallbackEngine == null) {
+                    singleDefaultContextKeyCallbackEngine = new 
DefaultContextKeyCallbackEngine();
+                    DefaultContextListenerManager 
instanceContextListenerManager = DefaultContextListenerManager.getInstance();
+                    
instanceContextListenerManager.getContextAsyncListenerBus().addListener(singleDefaultContextKeyCallbackEngine);
+                    logger.info("add listerner 
singleDefaultContextKeyCallbackEngine success");
+                }
+            }
+        }
+        return singleDefaultContextKeyCallbackEngine;
+    }
+
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/conf/ContextListenerConf.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/conf/ContextListenerConf.java
new file mode 100644
index 0000000000..808eb09491
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/conf/ContextListenerConf.java
@@ -0,0 +1,13 @@
+package com.webank.wedatasphere.linkis.cs.listener.conf;
+
+import com.webank.wedatasphere.linkis.common.conf.CommonVars;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/25
+ */
+public class ContextListenerConf {
+    public final static Integer WDS_CS_LISTENER_ASYN_CONSUMER_THREAD_MAX = 
Integer.parseInt(CommonVars.apply("wds.linkis.cs.listener.asyn.consumer.thread.max","5").getValue());
+    public final static Long 
WDS_CS_LISTENER_ASYN_CONSUMER_THREAD_FREE_TIME_MAX = 
Long.parseLong(CommonVars.apply("wds.linkis.cs.listener.asyn.consumer.freeTime.max","5000").getValue());
+    public final static Integer WDS_CS_LISTENER_ASYN_QUEUE_CAPACITY 
=Integer.parseInt(CommonVars.apply("wds.linkis.cs.listener.asyn.queue.size.max","300").getValue())
 ;
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextIDEvent.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextIDEvent.java
new file mode 100644
index 0000000000..2730744037
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextIDEvent.java
@@ -0,0 +1,13 @@
+package com.webank.wedatasphere.linkis.cs.listener.event;
+
+import com.webank.wedatasphere.linkis.common.listener.Event;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+
+/**
+ * @author peacewong
+ * @date 2020/2/20 15:03
+ */
+public interface ContextIDEvent extends Event {
+
+    ContextID getContextID();
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextKeyEvent.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextKeyEvent.java
new file mode 100644
index 0000000000..7cb1e6220b
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/ContextKeyEvent.java
@@ -0,0 +1,10 @@
+package com.webank.wedatasphere.linkis.cs.listener.event;
+
+import com.webank.wedatasphere.linkis.common.listener.Event;
+
+/**
+ * @author peacewong
+ * @date 2020/2/20 15:01
+ */
+public interface ContextKeyEvent extends Event {
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/enumeration/OperateType.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/enumeration/OperateType.java
new file mode 100644
index 0000000000..27b535865b
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/enumeration/OperateType.java
@@ -0,0 +1,13 @@
+package com.webank.wedatasphere.linkis.cs.listener.event.enumeration;
+
+/**
+ * @author peacewong
+ * @date 2020/2/15 21:33
+ */
+public enum OperateType {
+
+    /**
+     * 对contextKey和contextId的一些操作
+     */
+    ADD, UPDATE, DELETE, REMOVEALL, ACCESS, CREATE,REMOVE
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextIDEvent.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextIDEvent.java
new file mode 100644
index 0000000000..66c7cdc424
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextIDEvent.java
@@ -0,0 +1,34 @@
+package com.webank.wedatasphere.linkis.cs.listener.event.impl;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextIDEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.enumeration.OperateType;
+
+/**
+ * @author peacewong
+ * @date 2020/2/20 15:03
+ */
+public class DefaultContextIDEvent implements ContextIDEvent {
+
+    private ContextID contextID;
+
+    public OperateType getOperateType() {
+        return operateType;
+    }
+
+    public void setOperateType(OperateType operateType) {
+        this.operateType = operateType;
+    }
+
+    //TODO
+    private OperateType operateType;
+
+    @Override
+    public ContextID getContextID() {
+        return contextID;
+    }
+
+    public void setContextID(ContextID contextID) {
+        this.contextID = contextID;
+    }
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextKeyEvent.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextKeyEvent.java
new file mode 100644
index 0000000000..d308a01122
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/event/impl/DefaultContextKeyEvent.java
@@ -0,0 +1,56 @@
+package com.webank.wedatasphere.linkis.cs.listener.event.impl;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.enumeration.ContextType;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKeyValue;
+import com.webank.wedatasphere.linkis.cs.listener.event.ContextKeyEvent;
+import 
com.webank.wedatasphere.linkis.cs.listener.event.enumeration.OperateType;
+
+/**
+ * @author peacewong
+ * @date 2020/2/15 21:30
+ */
+public class DefaultContextKeyEvent implements ContextKeyEvent {
+
+    private ContextID contextID;
+
+    private ContextKeyValue contextKeyValue;
+
+    private ContextKeyValue oldValue;
+
+
+    private OperateType operateType;
+
+    public ContextID getContextID() {
+        return contextID;
+    }
+
+    public void setContextID(ContextID contextID) {
+        this.contextID = contextID;
+    }
+
+    public ContextKeyValue getContextKeyValue() {
+        return contextKeyValue;
+    }
+
+    public void setContextKeyValue(ContextKeyValue contextKeyValue) {
+        this.contextKeyValue = contextKeyValue;
+    }
+
+    public OperateType getOperateType() {
+        return operateType;
+    }
+
+    public void setOperateType(OperateType operateType) {
+        this.operateType = operateType;
+    }
+
+    public ContextKeyValue getOldValue() {
+        return oldValue;
+    }
+
+    public void setOldValue(ContextKeyValue oldValue) {
+        this.oldValue = oldValue;
+    }
+
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/ListenerManager.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/ListenerManager.java
new file mode 100644
index 0000000000..595ab80bbf
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/ListenerManager.java
@@ -0,0 +1,15 @@
+package com.webank.wedatasphere.linkis.cs.listener.manager;
+
+import 
com.webank.wedatasphere.linkis.cs.listener.ListenerBus.ContextAsyncListenerBus;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextIDCallbackEngine;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextKeyCallbackEngine;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/21
+ */
+public interface ListenerManager {
+     public ContextAsyncListenerBus getContextAsyncListenerBus(); //单例
+     public DefaultContextIDCallbackEngine getContextIDCallbackEngine(); //单例
+     public DefaultContextKeyCallbackEngine getContextKeyCallbackEngine(); //单例
+}
diff --git 
a/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/imp/DefaultContextListenerManager.java
 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/imp/DefaultContextListenerManager.java
new file mode 100644
index 0000000000..0dc473a413
--- /dev/null
+++ 
b/contextservice/cs-listener/src/main/java/com/webank/wedatasphere/linkis/cs/listener/manager/imp/DefaultContextListenerManager.java
@@ -0,0 +1,46 @@
+package com.webank.wedatasphere.linkis.cs.listener.manager.imp;
+
+import 
com.webank.wedatasphere.linkis.cs.listener.ListenerBus.ContextAsyncListenerBus;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextIDCallbackEngine;
+import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextKeyCallbackEngine;
+import com.webank.wedatasphere.linkis.cs.listener.manager.ListenerManager;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/21
+ */
+public class DefaultContextListenerManager implements ListenerManager {
+    @Override
+    public ContextAsyncListenerBus getContextAsyncListenerBus() {
+        ContextAsyncListenerBus contextAsyncListenerBus = 
ContextAsyncListenerBus.getInstance();
+        return contextAsyncListenerBus;
+    }
+
+    @Override
+    public DefaultContextIDCallbackEngine getContextIDCallbackEngine() {
+        DefaultContextIDCallbackEngine instanceIdCallbackEngine = 
DefaultContextIDCallbackEngine.getInstance();
+        return instanceIdCallbackEngine;
+    }
+
+    @Override
+    public DefaultContextKeyCallbackEngine getContextKeyCallbackEngine() {
+        DefaultContextKeyCallbackEngine instanceKeyCallbackEngine = 
DefaultContextKeyCallbackEngine.getInstance();
+        return instanceKeyCallbackEngine;
+    }
+
+    private static DefaultContextListenerManager 
singleDefaultContextListenerManager = null;
+
+    private DefaultContextListenerManager() {
+    }
+
+    public static DefaultContextListenerManager getInstance() {
+        if (singleDefaultContextListenerManager == null) {
+            synchronized (DefaultContextListenerManager.class) {
+                if (singleDefaultContextListenerManager == null) {
+                    singleDefaultContextListenerManager = new 
DefaultContextListenerManager();
+                }
+            }
+        }
+        return singleDefaultContextListenerManager;
+    }
+}
diff --git 
a/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextID.java
 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextID.java
new file mode 100644
index 0000000000..d82c0f5008
--- /dev/null
+++ 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextID.java
@@ -0,0 +1,22 @@
+package com.webank.wedatasphere.linkis.cs.listener.test;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+
+/**
+ * @author peacewong
+ * @date 2020/2/13 20:41
+ */
+public class TestContextID implements ContextID {
+
+    String contextID;
+
+    @Override
+    public String getContextId() {
+        return contextID;
+    }
+
+    @Override
+    public void setContextId(String contextId) {
+        this.contextID = contextId;
+    }
+}
diff --git 
a/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKey.java
 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKey.java
new file mode 100644
index 0000000000..2e0017cf68
--- /dev/null
+++ 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKey.java
@@ -0,0 +1,63 @@
+package com.webank.wedatasphere.linkis.cs.listener.test;
+
+import 
com.webank.wedatasphere.linkis.cs.common.entity.enumeration.ContextScope;
+import com.webank.wedatasphere.linkis.cs.common.entity.enumeration.ContextType;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/22
+ */
+public class TestContextKey implements ContextKey {
+    private  String key;
+    private ContextType contextType;
+    @Override
+    public String getKey() {
+        return this.key;
+    }
+
+    @Override
+    public void setKey(String key) {
+        this.key=key;
+    }
+
+    @Override
+    public ContextType getContextType() {
+        return this.contextType;
+    }
+
+    @Override
+    public void setContextType(ContextType contextType) {
+        this.contextType=contextType;
+    }
+
+    @Override
+    public ContextScope getContextScope() {
+        return null;
+    }
+
+    @Override
+    public void setContextScope(ContextScope contextScope) {
+
+    }
+
+    @Override
+    public String getKeywords() {
+        return null;
+    }
+
+    @Override
+    public void setKeywords(String keywords) {
+
+    }
+
+    @Override
+    public int getType() {
+        return 0;
+    }
+
+    @Override
+    public void setType(int type) {
+
+    }
+}
diff --git 
a/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKeyValue.java
 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKeyValue.java
new file mode 100644
index 0000000000..e8ed25310b
--- /dev/null
+++ 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextKeyValue.java
@@ -0,0 +1,36 @@
+package com.webank.wedatasphere.linkis.cs.listener.test;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKeyValue;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextValue;
+
+/**
+ * @author chaogefeng
+ * @date 2020/2/22 16:46
+ */
+public class TestContextKeyValue implements ContextKeyValue {
+
+    private ContextKey contextKey;
+
+    private ContextValue contextValue;
+
+    @Override
+    public ContextKey getContextKey() {
+        return this.contextKey;
+    }
+
+    @Override
+    public void setContextKey(ContextKey contextKey) {
+        this.contextKey = contextKey;
+    }
+
+    @Override
+    public ContextValue getContextValue() {
+        return this.contextValue;
+    }
+
+    @Override
+    public void setContextValue(ContextValue contextValue) {
+        this.contextValue = contextValue;
+    }
+}
diff --git 
a/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextValue.java
 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextValue.java
new file mode 100644
index 0000000000..41bc8ad820
--- /dev/null
+++ 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestContextValue.java
@@ -0,0 +1,37 @@
+package com.webank.wedatasphere.linkis.cs.listener.test;
+
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextValue;
+import com.webank.wedatasphere.linkis.cs.common.entity.source.ValueBean;
+
+/**
+ * @Author: chaogefeng
+ * @Date: 2020/2/22
+ */
+public class TestContextValue implements ContextValue {
+    private  Object value;
+
+    private String keywords;
+
+
+    @Override
+    public String getKeywords() {
+        return null;
+    }
+
+    @Override
+    public void setKeywords(String keywords) {
+
+    }
+
+    @Override
+    public Object getValue() {
+        return this.value;
+    }
+
+    @Override
+    public void setValue(Object value) {
+        this.value=value;
+    }
+
+
+}
diff --git 
a/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestListenerManager.java
 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestListenerManager.java
new file mode 100644
index 0000000000..4c44722c7b
--- /dev/null
+++ 
b/contextservice/cs-listener/src/test/java/com/webank/wedatasphere/linkis/cs/listener/test/TestListenerManager.java
@@ -0,0 +1,108 @@
+//package com.webank.wedatasphere.linkis.cs.listener.test;
+//
+//import com.webank.wedatasphere.linkis.common.listener.Event;
+//import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextID;
+//import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKey;
+//import 
com.webank.wedatasphere.linkis.cs.common.entity.source.ContextKeyValue;
+//import com.webank.wedatasphere.linkis.cs.common.entity.source.ContextValue;
+//import 
com.webank.wedatasphere.linkis.cs.listener.ListenerBus.ContextAsyncListenerBus;
+//import com.webank.wedatasphere.linkis.cs.listener.callback.imp.ClientSource;
+//import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.ContextKeyValueBean;
+//import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextIDCallbackEngine;
+//import 
com.webank.wedatasphere.linkis.cs.listener.callback.imp.DefaultContextKeyCallbackEngine;
+//import 
com.webank.wedatasphere.linkis.cs.listener.event.enumeration.OperateType;
+//import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextIDEvent;
+//import 
com.webank.wedatasphere.linkis.cs.listener.event.impl.DefaultContextKeyEvent;
+//import com.webank.wedatasphere.linkis.cs.listener.manager.ListenerManager;
+//import 
com.webank.wedatasphere.linkis.cs.listener.manager.imp.DefaultContextListenerManager;
+//import org.junit.Test;
+//
+//import java.util.ArrayList;
+//import java.util.List;
+//
+///**
+// * @Author: chaogefeng
+// * @Date: 2020/2/22
+// */
+//public class TestListenerManager {
+//    @Test
+//    public void testGetContextAsyncListenerBus() {
+//        DefaultContextListenerManager defaultContextListenerManager = 
DefaultContextListenerManager.getInstance();
+//
+//        ContextAsyncListenerBus contextAsyncListenerBus = 
defaultContextListenerManager.getContextAsyncListenerBus();
+//
+//        DefaultContextIDCallbackEngine contextIDCallbackEngine = 
defaultContextListenerManager.getContextIDCallbackEngine();
+//
+//        DefaultContextKeyCallbackEngine contextKeyCallbackEngine = 
defaultContextListenerManager.getContextKeyCallbackEngine();
+//        //client1的contextID
+//        TestContextID testContextID1 = new TestContextID();
+//        testContextID1.setContextId("18392881376");
+//
+//        //client2的contextID
+//        TestContextID testContextID2 = new TestContextID();
+//        testContextID2.setContextId("13431335441");
+//
+//        //client1的cskeys,监听key1,key2
+//        List<ContextKey> csKeys1 = new ArrayList<>();
+//        TestContextKey testContextKey1 = new TestContextKey();
+//        testContextKey1.setKey("key1");
+//        TestContextKey testContextKey2 = new TestContextKey();
+//        testContextKey2.setKey("key2");
+//        csKeys1.add(testContextKey1);
+//        csKeys1.add(testContextKey2);
+//
+//        //client2的cskeys,监听key3,key4
+//        List<ContextKey> csKeys2 = new ArrayList<>();
+//        TestContextKey testContextKey3 = new TestContextKey();
+//        testContextKey3.setKey("key3");
+//        TestContextKey testContextKey4 = new TestContextKey();
+//        testContextKey4.setKey("key4");
+//        csKeys2.add(testContextKey3);
+//        csKeys2.add(testContextKey4);
+//
+//
+//        //client1的 name及instance
+//        ClientSource clientSource1 = new ClientSource();
+//
+//        //client2的 name及instance
+//        ClientSource clientSource2 = new ClientSource();
+
+//
+//        contextKeyCallbackEngine.registerClient(testContextID1, csKeys1, 
clientSource1);
+//        
System.out.println("+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++");
+//
+//        contextKeyCallbackEngine.registerClient(testContextID2, csKeys2, 
clientSource2);
+//        //同一个client可监听多个contextID,同一个contextID 可以监听多个cskey
+//        contextKeyCallbackEngine.registerClient(testContextID1, csKeys1, 
clientSource2);
+//
+//
+//        //传递的事件,赋值 contextId:183292881376 
,ContextKeyValue:key1:chaogefeng,operateType的类型
+//        DefaultContextKeyEvent defaultContextKeyEvent = new 
DefaultContextKeyEvent();
+//        //1、设置ID
+//        defaultContextKeyEvent.setContextID(testContextID1);
+//        //2、设置操作类型
+//        defaultContextKeyEvent.setOperateType(OperateType.UPDATE);
+//        TestContextKeyValue testContextKeyValue = new TestContextKeyValue();
+//        //更新key1对应的值,修改为chaogefeng
+//        testContextKeyValue.setContextKey(testContextKey1);
+//        TestContextValue testContextValue = new TestContextValue();
+//        testContextValue.setValue("chaogefeng");
+//        testContextKeyValue.setContextValue(testContextValue);
+//        //3、设置key value值
+//        defaultContextKeyEvent.setContextKeyValue(testContextKeyValue);
+//        //4、给listener contextKeyCallbackEngine投递事件
+//        contextAsyncListenerBus.doPostEvent(contextKeyCallbackEngine, 
defaultContextKeyEvent);
+//        //5、cleint2 心跳,因为它也监听了 contextID1,则应该返回 contextID1 中最新的key value值给它
+//        ArrayList<ContextKeyValueBean> clientSource2ListenerCallback = 
contextKeyCallbackEngine.getListenerCallback(clientSource2);
+//        //遍历 这个bean,检查是否更新。
+//        
System.out.println("----------------------------------------------------------------------");
+//        for (ContextKeyValueBean contextKeyValueBean : 
clientSource2ListenerCallback) {
+//            System.out.println("返回的bean里面对应的contexID: " + 
contextKeyValueBean.getCsID().getContextId());
+//            System.out.println("返回的bean里面对应的cskeys: " + 
contextKeyValueBean.getCsKey().getKey());
+//            if (contextKeyValueBean.getCsValue() != null) {
+//                System.out.println("返回的bean里面对应的value: " + 
contextKeyValueBean.getCsValue().getValue());
+//            }
+//        }
+//    }
+//
+//}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to