This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new c6fb3de061 [INLONG-8163][Manager][DataProxy] Make DataProxy config
interface compatible with old versions (#8164)
c6fb3de061 is described below
commit c6fb3de0618736f4a633859820674de7ce0aa96e
Author: vernedeng <[email protected]>
AuthorDate: Tue Jun 6 21:00:29 2023 +0800
[INLONG-8163][Manager][DataProxy] Make DataProxy config interface
compatible with old versions (#8164)
---
.../dataproxy/sink/mq/pulsar/PulsarHandler.java | 8 ++---
.../service/cluster/InlongClusterService.java | 13 ++++++++
.../service/cluster/InlongClusterServiceImpl.java | 37 ++++++++++++++++++++++
.../repository/DataProxyConfigRepository.java | 10 ++++++
...itory.java => DataProxyConfigRepositoryV2.java} | 11 ++++---
.../controller/openapi/DataProxyController.java | 11 ++++++-
6 files changed, 80 insertions(+), 10 deletions(-)
diff --git
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
index 1b4d273d11..1800c5550a 100644
---
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
+++
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
@@ -66,7 +66,7 @@ public class PulsarHandler implements MessageQueueHandler {
// log print count
private static final LogCounter logCounter = new LogCounter(10, 100000, 30
* 1000);
- public static final String KEY_TENANT = "pulsarTenant";
+ public static final String KEY_TENANT = "tenant";
public static final String KEY_NAMESPACE = "namespace";
public static final String KEY_SERVICE_URL = "serviceUrl";
@@ -92,7 +92,7 @@ public class PulsarHandler implements MessageQueueHandler {
private String clusterName;
private MessageQueueZoneSinkContext sinkContext;
- private String pulsarTenant;
+ private String tenant;
private String namespace;
private ThreadLocal<EventHandler> handlerLocal = new ThreadLocal<>();
@@ -113,7 +113,7 @@ public class PulsarHandler implements MessageQueueHandler {
this.config = config;
this.clusterName = config.getClusterName();
this.sinkContext = sinkContext;
- this.pulsarTenant = config.getParams().get(KEY_TENANT);
+ this.tenant = config.getParams().get(KEY_TENANT);
this.namespace = config.getParams().get(KEY_NAMESPACE);
}
@@ -214,7 +214,7 @@ public class PulsarHandler implements MessageQueueHandler {
return false;
}
// topic
- String producerTopic = idConfig.getPulsarTopicName(pulsarTenant,
namespace);
+ String producerTopic = idConfig.getPulsarTopicName(tenant,
namespace);
if (producerTopic == null) {
sinkContext.fileMetricEventInc(StatConstants.EVENT_SINK_NOTOPIC);
sinkContext.addSendResultMetric(event, clusterName,
event.getUid(), false, 0);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterService.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterService.java
index e7c258350d..3191b1e69e 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterService.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterService.java
@@ -416,10 +416,23 @@ public interface InlongClusterService {
/**
* Get data proxy cluster list by the given cluster name
*
+ * This method was deprecated since version 1.8.0,
+ * new method please see {@link InlongClusterService#getMetaConfig(String,
String)}
+ *
* @return data proxy config
*/
+ @Deprecated
String getAllConfig(String clusterName, String md5);
+ /**
+ * Get data proxy cluster list by the given cluster name.
+ *
+ * since version 1.8.0
+ *
+ * @return data proxy config
+ */
+ String getMetaConfig(String clusterName, String md5);
+
/**
* Get the MQ info by cluster tag for Audit
*
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterServiceImpl.java
index f3e3e30807..36012fafae 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/cluster/InlongClusterServiceImpl.java
@@ -66,6 +66,7 @@ import org.apache.inlong.manager.pojo.user.UserInfo;
import
org.apache.inlong.manager.service.cluster.node.InlongClusterNodeOperator;
import
org.apache.inlong.manager.service.cluster.node.InlongClusterNodeOperatorFactory;
import org.apache.inlong.manager.service.repository.DataProxyConfigRepository;
+import
org.apache.inlong.manager.service.repository.DataProxyConfigRepositoryV2;
import org.apache.inlong.manager.service.user.UserService;
import com.github.pagehelper.Page;
@@ -120,10 +121,16 @@ public class InlongClusterServiceImpl implements
InlongClusterService {
private InlongClusterEntityMapper clusterMapper;
@Autowired
private InlongClusterNodeEntityMapper clusterNodeMapper;
+
@Lazy
@Autowired
+ @Deprecated
private DataProxyConfigRepository proxyRepository;
+ @Lazy
+ @Autowired
+ private DataProxyConfigRepositoryV2 proxyRepositoryV2;
+
@Override
public Integer saveTag(ClusterTagRequest request, String operator) {
LOGGER.debug("begin to save cluster tag {}", request);
@@ -1337,6 +1344,7 @@ public class InlongClusterServiceImpl implements
InlongClusterService {
}
@Override
+ @Deprecated
public String getAllConfig(String clusterName, String md5) {
DataProxyConfigResponse response = new DataProxyConfigResponse();
String configMd5 = proxyRepository.getProxyMd5(clusterName);
@@ -1365,6 +1373,35 @@ public class InlongClusterServiceImpl implements
InlongClusterService {
return configJson;
}
+ @Override
+ public String getMetaConfig(String clusterName, String md5) {
+ DataProxyConfigResponse response = new DataProxyConfigResponse();
+ String configMd5 = proxyRepositoryV2.getProxyMd5(clusterName);
+ if (configMd5 == null) {
+ response.setResult(false);
+ response.setErrCode(DataProxyConfigResponse.REQ_PARAMS_ERROR);
+ return GSON.toJson(response);
+ }
+
+ // same config
+ if (configMd5.equals(md5)) {
+ response.setResult(true);
+ response.setErrCode(DataProxyConfigResponse.NOUPDATE);
+ response.setMd5(configMd5);
+ response.setData(new DataProxyCluster());
+ return GSON.toJson(response);
+ }
+
+ String configJson = proxyRepositoryV2.getProxyConfigJson(clusterName);
+ if (configJson == null) {
+ response.setResult(false);
+ response.setErrCode(DataProxyConfigResponse.REQ_PARAMS_ERROR);
+ return GSON.toJson(response);
+ }
+
+ return configJson;
+ }
+
@Override
public AuditConfig getAuditConfig(String clusterTag) {
AuditConfig auditConfig = new AuditConfig();
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
index bbd974e525..263b91c52e 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
@@ -77,14 +77,18 @@ import java.util.concurrent.ConcurrentHashMap;
/**
* DataProxyConfigRepository
+ * This repository was deprecated since version 1.8.0
*/
@Lazy
+@Deprecated
@Repository(value = "dataProxyConfigRepository")
public class DataProxyConfigRepository implements IRepository {
public static final Logger LOGGER =
LoggerFactory.getLogger(DataProxyConfigRepository.class);
public static final String KEY_NAMESPACE = "namespace";
+ public static final String KEY_NEW_TENANT_KEY = "pulsarTenant";
+ public static final String KEY_OLD_TENANT_KEY = "tenant";
public static final String KEY_BACKUP_CLUSTER_TAG = "backup_cluster_tag";
public static final String KEY_BACKUP_TOPIC = "backup_topic";
public static final String KEY_SORT_TASK_NAME = "defaultSortTaskName";
@@ -317,6 +321,12 @@ public class DataProxyConfigRepository implements
IRepository {
} catch (Exception e) {
LOGGER.error("parse json string to map error", e);
}
+
+ // to be compatible with multi-tenancy #7914
+ String tenant = mapObj.get(KEY_NEW_TENANT_KEY);
+ mapObj.remove(KEY_NEW_TENANT_KEY);
+ mapObj.put(KEY_OLD_TENANT_KEY, tenant);
+
return mapObj;
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
similarity index 99%
copy from
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
copy to
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
index bbd974e525..a49dbdd4b1 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepository.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/repository/DataProxyConfigRepositoryV2.java
@@ -76,13 +76,14 @@ import java.util.TimerTask;
import java.util.concurrent.ConcurrentHashMap;
/**
- * DataProxyConfigRepository
+ * DataProxyConfigRepositoryV2
+ * Since version 1.8.0
*/
@Lazy
-@Repository(value = "dataProxyConfigRepository")
-public class DataProxyConfigRepository implements IRepository {
+@Repository(value = "dataProxyConfigRepositoryV2")
+public class DataProxyConfigRepositoryV2 implements IRepository {
- public static final Logger LOGGER =
LoggerFactory.getLogger(DataProxyConfigRepository.class);
+ public static final Logger LOGGER =
LoggerFactory.getLogger(DataProxyConfigRepositoryV2.class);
public static final String KEY_NAMESPACE = "namespace";
public static final String KEY_BACKUP_CLUSTER_TAG = "backup_cluster_tag";
@@ -440,7 +441,7 @@ public class DataProxyConfigRepository implements
IRepository {
*/
private void setReloadTimer() {
Timer reloadTimer = new Timer(true);
- TimerTask task = new
RepositoryTimerTask<DataProxyConfigRepository>(this);
+ TimerTask task = new
RepositoryTimerTask<DataProxyConfigRepositoryV2>(this);
reloadTimer.scheduleAtFixedRate(task, reloadInterval, reloadInterval);
}
diff --git
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/openapi/DataProxyController.java
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/openapi/DataProxyController.java
index 17de1fcac7..cea54c8c82 100644
---
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/openapi/DataProxyController.java
+++
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/openapi/DataProxyController.java
@@ -70,11 +70,20 @@ public class DataProxyController {
}
@PostMapping("/dataproxy/getAllConfig")
- @ApiOperation(value = "Get all proxy config")
+ @ApiOperation(value = "Get all proxy config. " +
+ "This method was deprecated since version 1.8.0. " +
+ "Please use new method /dataproxy/getMetaConfig")
+ @Deprecated
public String getAllConfig(@RequestBody DataProxyConfigRequest request) {
return clusterService.getAllConfig(request.getClusterName(),
request.getMd5());
}
+ @PostMapping("/dataproxy/getMetaConfig")
+ @ApiOperation(value = "Get all DataProxy meta config")
+ public String getMetaConfig(@RequestBody DataProxyConfigRequest request) {
+ return clusterService.getMetaConfig(request.getClusterName(),
request.getMd5());
+ }
+
@RequestMapping(value = "/changeClusterTag", method = RequestMethod.PUT)
@ApiOperation(value = "Change cluster tag and topic for inlong group id")
public Response<String> changeClusterTag(@RequestParam String
inlongGroupId, @RequestParam String clusterTag,