|
|
package com.yoho.datasync.producer.canal;
|
|
|
|
|
|
import com.sun.deploy.util.StringUtils;
|
|
|
import com.yoho.datasync.core.base.message.TableConfig;
|
|
|
import com.yoho.datasync.core.base.message.TableConfigLoader;
|
|
|
import lombok.Data;
|
|
|
import org.springframework.boot.context.properties.ConfigurationProperties;
|
|
|
import org.springframework.stereotype.Component;
|
|
|
|
|
|
import javax.annotation.PostConstruct;
|
|
|
import javax.annotation.Resource;
|
|
|
import java.util.ArrayList;
|
|
|
import java.util.List;
|
|
|
import java.util.stream.Collectors;
|
|
|
|
|
|
@Component
|
|
|
@ConfigurationProperties(prefix = "canal")
|
...
|
...
|
@@ -16,6 +23,9 @@ public class CanalConfig { |
|
|
|
|
|
private List<CanalInstance> canalInstance ;
|
|
|
|
|
|
@Resource
|
|
|
TableConfigLoader tableConfigLoader;
|
|
|
|
|
|
public List<CanalInstance> getCanalInstance() {
|
|
|
return canalInstance;
|
|
|
}
|
...
|
...
|
@@ -56,7 +66,27 @@ public class CanalConfig { |
|
|
|
|
|
String filter;
|
|
|
|
|
|
String dbname;
|
|
|
|
|
|
int fetchSize;
|
|
|
|
|
|
}
|
|
|
|
|
|
@PostConstruct
|
|
|
void buildInstanceFilter(){
|
|
|
canalInstance.forEach(instance -> {
|
|
|
String dbName = instance.getDbname();
|
|
|
String filter = buildCanalFilter(dbName);
|
|
|
instance.setFilter(filter);
|
|
|
});
|
|
|
}
|
|
|
|
|
|
private String buildCanalFilter(String dbName) {
|
|
|
List<TableConfig> dbTableConfigList = tableConfigLoader.getTableConfigs().stream().filter(a -> a.getDbName().equalsIgnoreCase(dbName)).collect(Collectors.toList());
|
|
|
List<String> filters = new ArrayList<>();
|
|
|
for (TableConfig tableConfig : dbTableConfigList) {
|
|
|
filters.add(dbName + "." + tableConfig.getTableName());
|
|
|
}
|
|
|
return StringUtils.join(filters, ",");
|
|
|
}
|
|
|
} |
...
|
...
|
|