001package gudusoft.gsqlparser.dlineage.dataflow.sqlenv; 002 003import gudusoft.gsqlparser.EDbVendor; 004import gudusoft.gsqlparser.dlineage.dataflow.model.SqlInfo; 005import gudusoft.gsqlparser.dlineage.metadata.Coordinate; 006import gudusoft.gsqlparser.sqlenv.*; 007import gudusoft.gsqlparser.sqlenv.parser.TJSONSQLEnvParser; 008import gudusoft.gsqlparser.sqlenv.parser.TSQLEnvParser; 009import gudusoft.gsqlparser.util.Logger; 010import gudusoft.gsqlparser.util.LoggerFactory; 011import gudusoft.gsqlparser.util.SQLUtil; 012 013import java.io.File; 014import java.util.*; 015import java.util.concurrent.CopyOnWriteArrayList; 016import java.util.concurrent.CountDownLatch; 017import java.util.concurrent.Executors; 018import java.util.concurrent.ThreadPoolExecutor; 019 020public class SQLEnvParser implements TSQLEnvParser { 021 022 private static final Logger logger = LoggerFactory.getLogger(Coordinate.class); 023 024 private TSQLEnv metadataSQLEnv; 025 026 private String defaultServer; 027 private String defaultDatabase; 028 private String defaultSchema; 029 private List<SqlInfo> metadataInfos = new CopyOnWriteArrayList<SqlInfo>(); 030 031 public SQLEnvParser(TSQLEnv metadataSQLEnv, String defaultServer, String defaultDatabase, String defaultSchema) { 032 this.metadataSQLEnv = metadataSQLEnv; 033 this.defaultServer = defaultServer; 034 this.defaultDatabase = defaultDatabase; 035 this.defaultSchema = defaultSchema; 036 } 037 038 public SQLEnvParser(String defaultServer, String defaultDatabase, String defaultSchema) { 039 this.metadataSQLEnv = null; 040 this.defaultServer = defaultServer; 041 this.defaultDatabase = defaultDatabase; 042 this.defaultSchema = defaultSchema; 043 } 044 045 public TSQLEnv[] parseSQLEnv(final EDbVendor vendor, SqlInfo[] sqlInfos) { 046 if (sqlInfos == null || sqlInfos.length == 0) { 047 return null; 048 } 049 final CountDownLatch latch = new CountDownLatch(sqlInfos.length); 050 final List<TSQLEnv> sqlenvs = new ArrayList<TSQLEnv>(); 051 int thread = Runtime.getRuntime().availableProcessors() / 4 + 1; 052 if (thread < 4) thread = 4; 053 ThreadPoolExecutor executor = (ThreadPoolExecutor) Executors 054 .newFixedThreadPool(thread < sqlInfos.length 055 ? thread 056 : sqlInfos.length); 057 for (int i = 0; i < sqlInfos.length; i++) { 058 final SqlInfo sqlInfo = sqlInfos[i]; 059 Runnable task = new Runnable() { 060 @Override 061 public void run() { 062 try { 063 if (sqlInfo == null) { 064 return; 065 } 066 String sql = sqlInfo.getSql(); 067 if (SQLUtil.isEmpty(sql) && sqlInfo.getFileName() != null 068 && new File(sqlInfo.getFileName()).exists()) { 069 sql = SQLUtil.getFileContent(sqlInfo.getFileName()); 070 } 071 if (SQLUtil.isEmpty(sql) && sqlInfo.getFilePath() != null 072 && new File(sqlInfo.getFilePath()).exists()) { 073 sql = SQLUtil.getFileContent(sqlInfo.getFilePath()); 074 } 075 TJSONSQLEnvParser jsonSQLEnvParser = new TJSONSQLEnvParser(defaultServer, defaultDatabase, defaultSchema); 076 // Pass the manifest's directory so a sharded manifest's relative 077 // catalog paths resolve against it, not the process working 078 // directory (a null baseDir could load a same-named catalog from 079 // CWD and merge a wrong environment). 080 String jsonBaseDir = sqlInfo.getFilePath() != null 081 ? new File(sqlInfo.getFilePath()).getParent() : null; 082 TSQLEnv[] sqlenv = jsonSQLEnvParser.parseSQLEnv(vendor, sql, jsonBaseDir); 083 if (sqlenv == null) { 084 sqlenv = parseSQLEnv(vendor, sql); 085 } 086 else { 087 metadataInfos.add(sqlInfo); 088 } 089 synchronized (sqlenvs) { 090 if (sqlenv != null) { 091 sqlenvs.addAll(Arrays.asList(sqlenv)); 092 } 093 } 094 } catch (Exception e) { 095 logger.error(e.getMessage(), e); 096 } finally { 097 latch.countDown(); 098 } 099 } 100 }; 101 executor.submit(task); 102 } 103 try { 104 latch.await(); 105 } catch (Exception e) { 106 logger.error("await latch occurs an exception.", e); 107 } 108 executor.shutdown(); 109 Map<String, List<TSQLEnv>> sqlenvMap = new HashMap<String, List<TSQLEnv>>(); 110 for (TSQLEnv sqlenv : sqlenvs) { 111 String serverName = sqlenv.getDefaultServerName(); 112 if(serverName == null) { 113 serverName= TSQLEnv.DEFAULT_SERVER_NAME; 114 } 115 if(!sqlenvMap.containsKey(serverName)) { 116 sqlenvMap.put(serverName, new ArrayList<TSQLEnv>()); 117 } 118 sqlenvMap.get(serverName).add(sqlenv); 119 } 120 121 List<TSQLEnv> mergeSQLEnvs = new ArrayList<TSQLEnv>(); 122 for(String key: sqlenvMap.keySet()) { 123 List<TSQLEnv> sqlenvItem = sqlenvMap.get(key); 124 mergeSQLEnvs.add(mergeSQLEnv(sqlenvItem)); 125 } 126 return mergeSQLEnvs.toArray(new TSQLEnv[0]); 127 } 128 129 public static TSQLEnv mergeSQLEnv(List<TSQLEnv> sqlenvs) { 130 if (sqlenvs.size() == 0) 131 return null; 132 TSQLEnv sqlenv = sqlenvs.get(0); 133 for (int i = 1; i < sqlenvs.size(); i++) { 134 TSQLEnv temp = sqlenvs.get(i); 135 if(SQLUtil.isEmpty(sqlenv.getDefaultServerName()) || sqlenv.getDefaultServerName().equals(TSQLEnv.DEFAULT_SERVER_NAME)) { 136 if(!SQLUtil.isEmpty(temp.getDefaultServerName()) && !temp.getDefaultServerName().equals(TSQLEnv.DEFAULT_SERVER_NAME)) { 137 sqlenv.setDefaultServerName(temp.getDefaultServerName()); 138 } 139 } 140 if (temp.getCatalogList() != null) { 141 for (int j = 0; j < temp.getCatalogList().size(); j++) { 142 mergeCatalog(sqlenv, temp.getCatalogList().get(j)); 143 } 144 } 145 } 146// if (sqlenv.getCatalogList() != null && sqlenv.getCatalogList().size() == 1) { 147// sqlenv.setDefaultCatalogName(sqlenv.getCatalogList().get(0).getName()); 148// if (sqlenv.getCatalogList().get(0).getSchemaList() != null 149// && sqlenv.getCatalogList().get(0).getSchemaList().size() == 1) { 150// sqlenv.setDefaultSchemaName(sqlenv.getCatalogList().get(0).getSchemaList().get(0).getName()); 151// } 152// } 153 return sqlenv; 154 } 155 156 private static void mergeCatalog(TSQLEnv sqlenv, TSQLCatalog catalog) { 157 TSQLCatalog mergeCatalog = sqlenv.getSQLCatalog(catalog.getName(), true); 158 List<TSQLSchema> schemaList = catalog.getSchemaList(); 159 if (schemaList != null) { 160 for (int i = 0; i < schemaList.size(); i++) { 161 mergeSchema(mergeCatalog, schemaList.get(i)); 162 } 163 } 164 } 165 166 private static void mergeSchema(TSQLCatalog catalog, TSQLSchema schema) { 167 TSQLSchema mergeSchema = catalog.getSchema(schema.getName(), true); 168 List<TSQLSchemaObject> schemaObjectList = schema.getSchemaObjectList(); 169 if (schemaObjectList != null) { 170 for (int i = 0; i < schemaObjectList.size(); i++) { 171 mergeSchemaObject(mergeSchema, schemaObjectList.get(i)); 172 } 173 } 174 } 175 176 private static void mergeSchemaObject(TSQLSchema schema, TSQLSchemaObject schemaObject) { 177 if (schemaObject instanceof TSQLTable) { 178 TSQLTable sqlTable = (TSQLTable) schemaObject; 179 TSQLTable mergeTable = schema.getSqlEnv().searchTable(schemaObject.getQualifiedName()); 180 if (mergeTable == null) { 181 mergeTable = schema.createTable(schemaObject.getName(), schemaObject.getPriority()); 182 mergeTable.setView(sqlTable.isView()); 183 } else if (schemaObject.getPriority() > mergeTable.getPriority()) { 184 mergeTable.setPriority(schemaObject.getPriority()); 185 } 186 List<TSQLColumn> columnList = sqlTable.getColumnList(); 187 if (columnList != null) { 188 for (int i = 0; i < columnList.size(); i++) { 189 mergeTable.addColumn(columnList.get(i).getName()); 190 } 191 } 192 } else if (schemaObject instanceof TSQLRoutine) { 193 TSQLRoutine sqlRoutine = (TSQLRoutine) schemaObject; 194 TSQLSchemaObject mergeRoutine = schema.getSqlEnv().searchSchemaObject(schemaObject.getQualifiedName(), sqlRoutine.getDataObjectType()); 195 if (mergeRoutine == null) { 196 mergeRoutine = schema.createSchemaObject(schemaObject.getName(), schemaObject.getDataObjectType()); 197 } 198 if (sqlRoutine.getDefinition() != null) { 199 ((TSQLRoutine) mergeRoutine).setDefinition(sqlRoutine.getDefinition()); 200 } 201 } else if (schemaObject instanceof TSQLSynonyms) { 202 TSQLSynonyms sqlSynonyms = (TSQLSynonyms) schemaObject; 203 TSQLSynonyms mergeSynonyms = schema.findSynonyms(schemaObject.getQualifiedName()); 204 if (mergeSynonyms == null) { 205 mergeSynonyms = schema.createSynonyms(schemaObject.getName(), 206 sqlSynonyms.getSourceDatabase(), 207 sqlSynonyms.getSourceSchema(), 208 sqlSynonyms.getSourceName()); 209 } else if (schemaObject.getPriority() > mergeSynonyms.getPriority()) { 210 mergeSynonyms.setPriority(schemaObject.getPriority()); 211 mergeSynonyms.setBaseTarget(sqlSynonyms.getSourceDatabase(), 212 sqlSynonyms.getSourceSchema(), 213 sqlSynonyms.getSourceName()); 214 } 215 } 216 } 217 218 public List<SqlInfo> getMetadataInfos() { 219 return metadataInfos; 220 } 221 222 public void setMetadataInfos(List<SqlInfo> metadataInfos) { 223 this.metadataInfos = metadataInfos; 224 } 225 226 @Override 227 public TSQLEnv[] parseSQLEnv(EDbVendor vendor, String sql) { 228 TDDLSQLEnv ddlSQLEnv = new TDDLSQLEnv(defaultServer, defaultDatabase, defaultSchema, metadataSQLEnv, vendor, sql); 229 ddlSQLEnv.initSQLEnv(); 230 if (ddlSQLEnv.isInit()) { 231 return new TSQLEnv[] { ddlSQLEnv }; 232 } else { 233 return null; 234 } 235 } 236}