add conflict file

This commit is contained in:
agapple
2019-08-31 12:47:27 +08:00
parent 4e0b8a0278
commit ba24c15833
8 changed files with 48 additions and 32 deletions
@@ -4,11 +4,10 @@ import com.alibaba.otter.canal.admin.common.exception.ServiceException;
import com.alibaba.otter.canal.protocol.exception.CanalClientException;
/**
* canal数据操作客户端
* canal admin操作客户端
*
* @author zebin.xuzb @ 2012-6-19
* @author jianghang
* @version 1.0.0
* @author agapple 2019年8月31日 下午12:46:29
* @since 1.1.4
*/
public interface AdminConnector {
@@ -3,12 +3,20 @@ package com.alibaba.otter.canal.admin.controller;
import java.util.List;
import java.util.Map;
import com.alibaba.otter.canal.admin.model.Pager;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import com.alibaba.otter.canal.admin.model.BaseModel;
import com.alibaba.otter.canal.admin.model.CanalInstanceConfig;
import com.alibaba.otter.canal.admin.model.Pager;
import com.alibaba.otter.canal.admin.service.CanalInstanceService;
/**
@@ -134,8 +142,7 @@ public class CanalInstanceController {
* @return 是否成功
*/
@PutMapping(value = "/instance/status/{id}")
public BaseModel<Boolean> instanceStart(@PathVariable Long id, @RequestParam String option,
@PathVariable String env) {
public BaseModel<Boolean> instanceStart(@PathVariable Long id, @RequestParam String option, @PathVariable String env) {
return BaseModel.getInstance(canalInstanceConfigService.instanceOperation(id, option));
}
@@ -162,6 +169,6 @@ public class CanalInstanceController {
*/
@GetMapping(value = "/active/instances/{serverId}")
public BaseModel<List<CanalInstanceConfig>> activeInstances(@PathVariable Long serverId, @PathVariable String env) {
return BaseModel.getInstance(canalInstanceConfigService.findActiveInstaceByServerId(serverId));
return BaseModel.getInstance(canalInstanceConfigService.findActiveInstanceByServerId(serverId));
}
}
@@ -20,8 +20,6 @@ public interface CanalInstanceService {
CanalInstanceConfig detail(Long id);
CanalInstanceConfig findOne(String name);
void updateContent(CanalInstanceConfig canalInstanceConfig);
void delete(Long id);
@@ -32,5 +30,5 @@ public interface CanalInstanceService {
boolean instanceOperation(Long id, String option);
List<CanalInstanceConfig> findActiveInstaceByServerId(Long serverId);
List<CanalInstanceConfig> findActiveInstanceByServerId(Long serverId);
}
@@ -8,6 +8,7 @@ import org.springframework.stereotype.Service;
import com.alibaba.otter.canal.admin.common.exception.ServiceException;
import com.alibaba.otter.canal.admin.model.CanalCluster;
import com.alibaba.otter.canal.admin.model.CanalInstanceConfig;
import com.alibaba.otter.canal.admin.model.NodeServer;
import com.alibaba.otter.canal.admin.service.CanalClusterServic;
@@ -33,6 +34,12 @@ public class CanalClusterServiceImpl implements CanalClusterServic {
throw new ServiceException("当前集群下存在Server, 无法删除");
}
// 判断集群下是否存在instance信息
int instanceCnt = CanalInstanceConfig.find.query().where().eq("clusterId", id).findCount();
if (instanceCnt > 0) {
throw new ServiceException("当前集群下存在Instance配置,无法删除");
}
CanalCluster canalCluster = CanalCluster.find.byId(id);
if (canalCluster != null) {
canalCluster.delete();
@@ -124,7 +124,7 @@ public class CanalInstanceServiceImpl implements CanalInstanceService {
*
* @param serverId server id
*/
public List<CanalInstanceConfig> findActiveInstaceByServerId(Long serverId) {
public List<CanalInstanceConfig> findActiveInstanceByServerId(Long serverId) {
NodeServer nodeServer = NodeServer.find.byId(serverId);
if (nodeServer == null) {
return null;
@@ -233,16 +233,6 @@ public class CanalInstanceServiceImpl implements CanalInstanceService {
}
}
@Override
public CanalInstanceConfig findOne(String name) {
CanalInstanceConfig config = CanalInstanceConfig.find.query()
.setDisableLazyLoading(true)
.where()
.eq("name", name)
.findOne();
return config;
}
public Map<String, String> remoteInstanceLog(Long id, Long nodeId) {
Map<String, String> result = new HashMap<>();
@@ -375,7 +375,7 @@ public class CanalController {
}
if (config.getMode().isManager()) {
PlainCanalInstanceGenerator instanceGenerator = new PlainCanalInstanceGenerator();
PlainCanalInstanceGenerator instanceGenerator = new PlainCanalInstanceGenerator(properties);
instanceGenerator.setCanalConfigClient(managerClients.get(config.getManagerAddress()));
instanceGenerator.setSpringXml(config.getSpringXml());
return instanceGenerator.generate(destination);
@@ -50,11 +50,8 @@ public class CanalLauncher {
if (StringUtils.isNotEmpty(managerAddress)) {
String user = properties.getProperty(CanalConstants.CANAL_ADMIN_USER);
String passwd = properties.getProperty(CanalConstants.CANAL_ADMIN_PASSWD);
String adminPort = properties.getProperty(CanalConstants.CANAL_ADMIN_PORT);
String adminPort = properties.getProperty(CanalConstants.CANAL_ADMIN_PORT, "11110");
String registerIp = properties.getProperty(CanalConstants.CANAL_REGISTER_IP);
if (StringUtils.isEmpty(adminPort)) {
adminPort = "11110";
}
if (StringUtils.isEmpty(registerIp)) {
registerIp = AddressUtils.getHostIp();
}
@@ -69,7 +66,9 @@ public class CanalLauncher {
+ " can't not found config for [" + registerIp + ":" + adminPort
+ "]");
}
properties = canalConfig.getProperties();
Properties managerProperties = canalConfig.getProperties();
// merge local
managerProperties.putAll(properties);
int scanIntervalInSecond = Integer.valueOf(properties.getProperty(CanalConstants.CANAL_AUTO_SCAN_INTERVAL,
"5"));
executor.scheduleWithFixedDelay(new Runnable() {
@@ -85,7 +84,10 @@ public class CanalLauncher {
if (newCanalConfig != null) {
// 远程配置canal.properties修改重新加载整个应用
canalStater.stop();
canalStater.setProperties(newCanalConfig.getProperties());
Properties managerProperties = newCanalConfig.getProperties();
// merge local
managerProperties.putAll(properties);
canalStater.setProperties(managerProperties);
canalStater.start();
lastCanalConfig = newCanalConfig;
@@ -98,9 +100,11 @@ public class CanalLauncher {
}
}, 0, scanIntervalInSecond, TimeUnit.SECONDS);
canalStater.setProperties(managerProperties);
} else {
canalStater.setProperties(properties);
}
canalStater.setProperties(properties);
canalStater.start();
runningLatch.await();
executor.shutdownNow();
@@ -1,5 +1,7 @@
package com.alibaba.otter.canal.instance.manager;
import java.util.Properties;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.BeanFactory;
@@ -26,6 +28,11 @@ public class PlainCanalInstanceGenerator implements CanalInstanceGenerator {
private PlainCanalConfigClient canalConfigClient;
private String defaultName = "instance";
private BeanFactory beanFactory;
private Properties canalConfig;
public PlainCanalInstanceGenerator(Properties canalConfig){
this.canalConfig = canalConfig;
}
public CanalInstance generate(String destination) {
synchronized (CanalInstanceGenerator.class) {
@@ -34,8 +41,12 @@ public class PlainCanalInstanceGenerator implements CanalInstanceGenerator {
if (canal == null) {
throw new CanalException("instance : " + destination + " config is not found");
}
Properties properties = canal.getProperties();
// merge local
properties.putAll(canalConfig);
// 设置动态properties,替换掉本地properties
com.alibaba.otter.canal.instance.spring.support.PropertyPlaceholderConfigurer.propertiesLocal.set(canal.getProperties());
com.alibaba.otter.canal.instance.spring.support.PropertyPlaceholderConfigurer.propertiesLocal.set(properties);
// 设置当前正在加载的通道,加载spring查找文件时会用到该变量
System.setProperty("canal.instance.destination", destination);
this.beanFactory = getBeanFactory(springXml);