001package gudusoft.gsqlparser.dlineage.metadata; 002 003import gudusoft.gsqlparser.EDbVendor; 004import gudusoft.gsqlparser.dlineage.dataflow.metadata.MetadataReader; 005import gudusoft.gsqlparser.dlineage.dataflow.metadata.sqlflow.SqlflowMetadataAnalyzer; 006import gudusoft.gsqlparser.dlineage.dataflow.metadata.sqlflow.sharded.SqlflowShardedMetadataAnalyzer; 007import gudusoft.gsqlparser.dlineage.dataflow.model.SubType; 008import gudusoft.gsqlparser.dlineage.dataflow.model.xml.*; 009import gudusoft.gsqlparser.dlineage.util.DataflowUtility; 010import gudusoft.gsqlparser.sqlenv.TSQLEnv; 011import gudusoft.gsqlparser.sqlenv.constant.SystemConstant; 012import gudusoft.gsqlparser.util.SQLUtil; 013import gudusoft.gsqlparser.util.json.JSON; 014 015import java.io.File; 016import java.util.ArrayList; 017import java.util.List; 018import java.util.Map; 019 020public class MetadataUtil { 021 022 public static Sqlflow convertMetadataJsonToSqlflow(EDbVendor vendor, String json) { 023 if (MetadataReader.isSqlflow(json)) { 024 dataflow temp = new SqlflowMetadataAnalyzer().analyzeMetadata(vendor, json); 025 Sqlflow sqlflow = MetadataUtil.convertDataflowToMetadata(vendor, temp); 026 return sqlflow; 027 } else { 028 throw new IllegalArgumentException("Not a sqlflow json."); 029 } 030 } 031 032 public static Sqlflow convertShardedMetadataJsonToSqlflow(EDbVendor vendor, String json, String baseDir) { 033 if (MetadataReader.isSqlflowSharded(json)) { 034 dataflow temp = new SqlflowShardedMetadataAnalyzer(baseDir).analyzeMetadata(vendor, json); 035 Sqlflow sqlflow = MetadataUtil.convertDataflowToMetadata(vendor, temp); 036 return sqlflow; 037 } else { 038 throw new IllegalArgumentException("Not a sqlflow-sharded json."); 039 } 040 } 041 042 public static Sqlflow convertDataflowToMetadata(EDbVendor vendor, dataflow dataflow) { 043 if (dataflow == null) { 044 return null; 045 } 046 Sqlflow sqlflow = new Sqlflow(); 047 sqlflow.setCreatedBy(SystemConstant.name + " " + SystemConstant.version); 048 appendTables(sqlflow, vendor, dataflow); 049 appendPackages(sqlflow, vendor, dataflow); 050 appendProcedures(sqlflow, vendor, dataflow); 051 appendProcesses(sqlflow, vendor, dataflow); 052 appendErrorMessages(sqlflow, vendor, dataflow); 053 return sqlflow; 054 } 055 056 private static void appendErrorMessages(Sqlflow sqlflow, EDbVendor vendor, dataflow dataflow) { 057 if (dataflow.getErrors() != null && !dataflow.getErrors().isEmpty()) { 058 for (error errorItem : dataflow.getErrors()) { 059 Error error = new Error(); 060 error.setCoordinate(errorItem.getCoordinate()); 061 error.setErrorMessage(errorItem.getErrorMessage()); 062 error.setErrorType(errorItem.getErrorType()); 063 error.setFile(errorItem.getFile()); 064 error.setOriginCoordinate(errorItem.getOriginCoordinate()); 065 sqlflow.appendError(error); 066 } 067 } 068 } 069 070 private static void appendProcedures(Sqlflow sqlflow, EDbVendor vendor, dataflow dataflow) { 071 if (dataflow.getProcedures() == null || dataflow.getProcedures().isEmpty()) { 072 return; 073 } 074 075 for (procedure procedureItem : dataflow.getProcedures()) { 076 Server server = new Server(); 077 server.setDbVendor(vendor.name()); 078 String serverName = procedureItem.getServer(); 079 if (serverName == null && !server.isSupportsCatalogs() && procedureItem.getDatabase() != null) { 080 serverName = procedureItem.getDatabase(); 081 } 082 if (serverName == null) { 083 serverName = TSQLEnv.DEFAULT_SERVER_NAME; 084 } 085 server.setName(serverName); 086 server = sqlflow.appendServer(server); 087 088 Database database = null; 089 090 if (server.isSupportsCatalogs()) { 091 String databaseName = procedureItem.getDatabase(); 092 database = new Database(); 093 database.setName(databaseName == null ? TSQLEnv.DEFAULT_DB_NAME : databaseName); 094 database = server.appendDatabase(database); 095 database.setServer(server); 096 } 097 098 Schema schema = null; 099 if (server.isSupportsSchemas()) { 100 String schemaName = procedureItem.getSchema(); 101 schema = new Schema(); 102 schema.setName(schemaName == null ? TSQLEnv.DEFAULT_SCHEMA_NAME : schemaName); 103 if (database != null) { 104 schema = database.appendSchema(server, schema); 105 schema.setParent(database); 106 schema.setServer(server); 107 } else { 108 schema = server.appendSchema(schema); 109 schema.setServer(server); 110 } 111 } 112 113 Procedure procedure = new Procedure(); 114 procedure.setId(procedureItem.getId()); 115 116 List<String> segments = SQLUtil.parseNames(procedureItem.getName()); 117 procedure.setName(segments.get(segments.size()-1)); 118 procedure.setServer(server); 119 if (server.isSupportsSchemas()) { 120 procedure.setParent(schema); 121 } 122 else if (server.isSupportsCatalogs()) { 123 procedure.setParent(database); 124 } 125 126 procedure.setType(procedureItem.getType()); 127 procedure.setCoordinates(Coordinate.parse(procedureItem.getCoordinate())); 128 if (procedureItem.getArguments() != null) { 129 for (argument argumentItem : procedureItem.getArguments()) { 130 Argument argument = new Argument(); 131 argument.setServer(server); 132 argument.setParent(procedure); 133 argument.setId(argumentItem.getId()); 134 argument.setName(argumentItem.getName()); 135 argument.setDataType(argumentItem.getDatatype()); 136 argument.setInout(argumentItem.getInout()); 137 argument.setCoordinates(Coordinate.parse(argumentItem.getCoordinate())); 138 procedure.appendArgument(argument); 139 } 140 } 141 142 if (server.isSupportsSchemas()) { 143 schema.appendProcedure(procedure); 144 } else if (server.isSupportsCatalogs()) { 145 database.appendProcedure(procedure); 146 } 147 } 148 } 149 150 protected static void appendPackages(Sqlflow sqlflow, EDbVendor vendor, dataflow dataflow) { 151 if (dataflow.getPackages() == null || dataflow.getPackages().isEmpty()) { 152 return; 153 } 154 for (oraclePackage oraclePackageItem : dataflow.getPackages()) { 155 Server server = new Server(); 156 server.setDbVendor(vendor.name()); 157 String serverName = oraclePackageItem.getServer(); 158 if (serverName == null && !server.isSupportsCatalogs() && oraclePackageItem.getDatabase() != null) { 159 serverName = oraclePackageItem.getDatabase(); 160 } 161 if (serverName == null) { 162 serverName = TSQLEnv.DEFAULT_SERVER_NAME; 163 } 164 server.setName(serverName); 165 server = sqlflow.appendServer(server); 166 167 Database database = null; 168 169 if (server.isSupportsCatalogs()) { 170 String databaseName = oraclePackageItem.getDatabase(); 171 database = new Database(); 172 database.setName(databaseName == null ? TSQLEnv.DEFAULT_DB_NAME : databaseName); 173 database = server.appendDatabase(database); 174 database.setServer(server); 175 } 176 177 Schema schema = null; 178 if (server.isSupportsSchemas()) { 179 String schemaName = oraclePackageItem.getSchema(); 180 schema = new Schema(); 181 schema.setName(schemaName == null ? TSQLEnv.DEFAULT_SCHEMA_NAME : schemaName); 182 if (database != null) { 183 schema = database.appendSchema(server, schema); 184 schema.setParent(database); 185 schema.setServer(server); 186 } else { 187 schema = server.appendSchema(schema); 188 schema.setServer(server); 189 } 190 } 191 192 Package oraclePackage = new Package(); 193 oraclePackage.setId(oraclePackageItem.getId()); 194 List<String> segments = SQLUtil.parseNames(oraclePackageItem.getName()); 195 oraclePackage.setName(segments.get(segments.size()-1)); 196 oraclePackage.setServer(server); 197 if (server.isSupportsSchemas()) { 198 oraclePackage.setParent(schema); 199 } 200 else if (server.isSupportsCatalogs()) { 201 oraclePackage.setParent(database); 202 } 203 204 oraclePackage.setCoordinates(Coordinate.parse(oraclePackageItem.getCoordinate())); 205 if (server.isSupportsSchemas()) { 206 schema.appendPackage(oraclePackage); 207 } 208 209 if (oraclePackageItem != null) { 210 for (procedure procedureItem : oraclePackageItem.getProcedures()) { 211 Procedure procedure = new Procedure(); 212 List<String> procedureSegments = SQLUtil.parseNames(procedureItem.getName()); 213 procedure.setName(procedureSegments.get(procedureSegments.size() - 1)); 214 procedure.setServer(server); 215 procedure.setParent(oraclePackage); 216 procedure.setId(procedureItem.getId()); 217 procedure.setType(procedureItem.getType()); 218 procedure.setCoordinates(Coordinate.parse(procedureItem.getCoordinate())); 219 if (procedureItem.getArguments() != null) { 220 for (argument argumentItem : procedureItem.getArguments()) { 221 Argument argument = new Argument(); 222 argument.setServer(server); 223 argument.setParent(procedure); 224 argument.setId(argumentItem.getId()); 225 argument.setName(argumentItem.getName()); 226 argument.setDataType(argumentItem.getDatatype()); 227 argument.setInout(argumentItem.getInout()); 228 argument.setCoordinates(Coordinate.parse(argumentItem.getCoordinate())); 229 procedure.appendArgument(argument); 230 } 231 } 232 oraclePackage.appendProcedure(procedure); 233 } 234 } 235 } 236 } 237 238 protected static void appendTables(Sqlflow sqlflow, EDbVendor vendor, dataflow dataflow) { 239 Map<String, table> tableMap = DataflowUtility.getDataflowDbObjMap(dataflow); 240 for (table tableItem : tableMap.values()) { 241 Server server = new Server(); 242 server.setDbVendor(vendor.name()); 243 String serverName = tableItem.getServer(); 244 if (serverName == null && !server.isSupportsCatalogs() && tableItem.getDatabase() != null) { 245 serverName = tableItem.getDatabase(); 246 } 247 if (serverName == null) { 248 serverName = TSQLEnv.DEFAULT_SERVER_NAME; 249 } 250 server.setName(serverName); 251 server = sqlflow.appendServer(server); 252 253 Database database = null; 254 255 if (server.isSupportsCatalogs()) { 256 String databaseName = tableItem.getDatabase(); 257 database = new Database(); 258 database.setName(databaseName == null ? TSQLEnv.DEFAULT_DB_NAME : databaseName); 259 database = server.appendDatabase(database); 260 database.setServer(server); 261 } 262 263 Schema schema = null; 264 if (server.isSupportsSchemas()) { 265 String schemaName = tableItem.getSchema(); 266 schema = new Schema(); 267 schema.setName(schemaName == null ? TSQLEnv.DEFAULT_SCHEMA_NAME : schemaName); 268 if (database != null) { 269 schema = database.appendSchema(server, schema); 270 schema.setParent(database); 271 schema.setServer(server); 272 } else { 273 schema = server.appendSchema(schema); 274 schema.setParent(schema); 275 schema.setServer(server); 276 } 277 } 278 279 Table table = new Table(); 280 table.setId(tableItem.getId()); 281 table.setType(tableItem.getType()); 282 if (tableItem.getAlias() != null) { 283 table.setAlias(tableItem.getAlias()); 284 } 285 table.setCoordinates(Coordinate.parse(tableItem.getCoordinate())); 286 table.setServer(server); 287 if (server.isSupportsSchemas()) { 288 table.setParent(schema); 289 } 290 else if (server.isSupportsCatalogs()) { 291 table.setParent(database); 292 } 293 table.setSubType(tableItem.getSubType()); 294 table.setFromDDL(tableItem.getFromDDL()); 295 table.setDisplayName(tableItem.getName()); 296 if (SubType.dblink.name().equals(tableItem.getSubType())) { 297 if (SQLUtil.trimColumnStringQuote(tableItem.getName().substring(tableItem.getName().lastIndexOf("@") + 1).trim()) 298 .equals(SQLUtil.trimColumnStringQuote(tableItem.getDatabase()))) { 299 List<String> segments = SQLUtil 300 .parseNames(tableItem.getName().substring(0, tableItem.getName().lastIndexOf("@")).trim()); 301 table.setName(segments.get(segments.size() - 1)); 302 table.setDbLink(tableItem.getDatabase()); 303 } else { 304 List<String> segments = SQLUtil.parseNames(tableItem.getName()); 305 table.setName(segments.get(segments.size() - 1)); 306 } 307 } else { 308 List<String> segments = SQLUtil.parseNames(tableItem.getName()); 309 table.setName(segments.get(segments.size() - 1)); 310 } 311 312 if (tableItem.getColumns() != null) { 313 List<Column> columns = new ArrayList<>(tableItem.getColumns().size()); 314 for (column columnItem : tableItem.getColumns()) { 315 Column column = new Column(); 316 column.setId(columnItem.getId()); 317 column.setName(columnItem.getName()); 318 column.setCoordinates(Coordinate.parse(columnItem.getCoordinate())); 319 column.setSource(columnItem.getSource()); 320 column.setParent(table); 321 column.setServer(server); 322 column.setDataType(columnItem.getDataType()); 323 column.setForeignKey(columnItem.isForeignKey()); 324 column.setPrimaryKey(columnItem.isPrimaryKey()); 325 column.setIndexKey(columnItem.isIndexKey()); 326 column.setUnqiueKey(columnItem.isUnqiueKey()); 327 column.setNullable(columnItem.isNullable()); 328 column.setIdentity(columnItem.isIdentity()); 329 column.setComputed(columnItem.isComputed()); 330 column.setPersisted(columnItem.isPersisted()); 331 column.setComputedDefinition(columnItem.getComputedDefinition()); 332 column.setDefaultValue(columnItem.getDefaultValue()); 333 column.setDefaultName(columnItem.getDefaultName()); 334 columns.add(column); 335 } 336 table.setColumns(columns); 337 } 338 339 if (server.isSupportsSchemas()) { 340 schema.appendTable(table); 341 } else if (server.isSupportsCatalogs()) { 342 database.appendTable(table); 343 } 344 } 345 } 346 347 protected static void appendProcesses(Sqlflow sqlflow, EDbVendor vendor, dataflow dataflow) { 348 if (dataflow.getProcesses() == null || dataflow.getProcesses().isEmpty()) { 349 return; 350 } 351 for (process processItem : dataflow.getProcesses()) { 352 Server server = new Server(); 353 server.setDbVendor(vendor.name()); 354 String serverName = processItem.getServer(); 355 if (serverName == null && !server.isSupportsCatalogs() && processItem.getDatabase() != null) { 356 serverName = processItem.getDatabase(); 357 } 358 if (serverName == null) { 359 serverName = TSQLEnv.DEFAULT_SERVER_NAME; 360 } 361 server.setName(serverName); 362 server = sqlflow.appendServer(server); 363 364 Database database = null; 365 366 if (server.isSupportsCatalogs()) { 367 String databaseName = processItem.getDatabase(); 368 database = new Database(); 369 database.setName(databaseName == null ? TSQLEnv.DEFAULT_DB_NAME : databaseName); 370 database = server.appendDatabase(database); 371 database.setServer(server); 372 } 373 374 Schema schema = null; 375 if (server.isSupportsSchemas()) { 376 String schemaName = processItem.getSchema(); 377 schema = new Schema(); 378 schema.setName(schemaName == null ? TSQLEnv.DEFAULT_SCHEMA_NAME : schemaName); 379 if (database != null) { 380 schema = database.appendSchema(server, schema); 381 schema.setParent(database); 382 schema.setServer(server); 383 } else { 384 schema = server.appendSchema(schema); 385 schema.setServer(server); 386 } 387 } 388 389 Process process = new Process(); 390 process.setId(processItem.getId()); 391 process.setName(processItem.getName()); 392 process.setProcedureId(processItem.getProcedureId()); 393 process.setProcedureName(processItem.getProcedureName()); 394 process.setQueryHashId(processItem.getQueryHashId()); 395 process.setCoordinates(Coordinate.parse(processItem.getCoordinate())); 396 process.setServer(server); 397 if (server.isSupportsSchemas()) { 398 process.setParent(schema); 399 } 400 else if (server.isSupportsCatalogs()) { 401 process.setParent(database); 402 } 403 if (processItem.getTransforms() != null && !processItem.getTransforms().isEmpty()) { 404 List<Transform> transforms = new ArrayList<Transform>(processItem.getTransforms().size()); 405 for(transform transformItem:processItem.getTransforms()) { 406 Transform transform = new Transform(); 407 transform.setCode(transformItem.getCode()); 408 transform.setType(transformItem.getType()); 409 transform.setCoordinate(transformItem.getCoordinate()); 410 transforms.add(transform); 411 } 412 process.setTransforms(transforms.toArray(new Transform[0])); 413 } 414 if (server.isSupportsSchemas()) { 415 schema.appendProcess(process); 416 } else if (server.isSupportsCatalogs()) { 417 database.appendProcess(process); 418 } 419 } 420 } 421 422 public static void main(String[] args) { 423// dataflow dataflow = XML2Model.loadXML(dataflow.class, 424// SQLUtil.getFileContent("C:\\Users\\KK\\Desktop\\dataflow.xml")); 425// Sqlflow sqlflow = MetadataUtil.convertDataflowToMetadata(EDbVendor.dbvoracle, dataflow); 426// System.out.println(JSON.toJSONString(sqlflow)); 427 428 String content = SQLUtil.getFileContent(new File("D:\\1.json")); 429 if (MetadataReader.isSqlflow(content)) { 430 dataflow temp = new SqlflowMetadataAnalyzer().analyzeMetadata(EDbVendor.dbvpostgresql, content); 431 Sqlflow sqlflow = MetadataUtil.convertDataflowToMetadata(EDbVendor.dbvpostgresql, temp); 432 System.out.println(JSON.toJSONString(sqlflow)); 433 } 434 } 435}