001package gudusoft.gsqlparser.dlineage.dataflow.model; 002 003import gudusoft.gsqlparser.EDbVendor; 004import gudusoft.gsqlparser.ESetOperatorType; 005import gudusoft.gsqlparser.TCustomSqlStatement; 006import gudusoft.gsqlparser.TSourceToken; 007import gudusoft.gsqlparser.dlineage.util.DlineageUtil; 008import gudusoft.gsqlparser.dlineage.util.Pair3; 009import gudusoft.gsqlparser.nodes.TCTE; 010import gudusoft.gsqlparser.nodes.TFunctionCall; 011import gudusoft.gsqlparser.nodes.TJoinList; 012import gudusoft.gsqlparser.pp.utils.SourceTokenSearcher; 013import gudusoft.gsqlparser.sqlenv.ESQLDataObjectType; 014import gudusoft.gsqlparser.sqlenv.TSQLEnv; 015import gudusoft.gsqlparser.stmt.*; 016import gudusoft.gsqlparser.stmt.hive.THiveLoad; 017import gudusoft.gsqlparser.stmt.mssql.TCreateExternalDataSourceStmt; 018import gudusoft.gsqlparser.stmt.mssql.TMssqlCreateFunction; 019import gudusoft.gsqlparser.stmt.mssql.TMssqlDeclare; 020import gudusoft.gsqlparser.stmt.mssql.TMssqlExecute; 021import gudusoft.gsqlparser.stmt.snowflake.TCreateStageStmt; 022import gudusoft.gsqlparser.stmt.snowflake.TCreateStreamStmt; 023import gudusoft.gsqlparser.stmt.snowflake.TSnowflakeCopyIntoStmt; 024import gudusoft.gsqlparser.stmt.teradata.TTeradataCreateProcedure; 025import gudusoft.gsqlparser.util.Logger; 026import gudusoft.gsqlparser.util.LoggerFactory; 027import gudusoft.gsqlparser.util.SQLUtil; 028 029import java.io.ByteArrayInputStream; 030import java.io.IOException; 031import java.util.ArrayDeque; 032import java.util.ArrayList; 033import java.util.Deque; 034import java.util.List; 035import java.util.Properties; 036 037public class Process { 038 private static final Logger logger = LoggerFactory.getLogger(Process.class); 039 private long id; 040 protected String server; 041 protected String schema; 042 protected String database; 043 private String name; 044 private String procedureName; 045 private Long procedureId; 046 private String queryHashId; 047 private String type; 048 private String customType; 049 private TCustomSqlStatement gspObject; 050 private Pair3<Long, Long, String> startPosition; 051 private Pair3<Long, Long, String> endPosition; 052 private List<Object> targetColumns = new ArrayList<Object>(); 053 private List<Transform> transforms = new ArrayList<Transform>(); 054 055 public Process(TCustomSqlStatement gspObject) { 056 057 if (gspObject == null) { 058 throw new IllegalArgumentException("Process arguments can't be null."); 059 } 060 061 id = ++ModelBindingManager.get().TABLE_COLUMN_ID; 062 063 this.gspObject = gspObject; 064 065 TSourceToken startToken = gspObject.getStartToken(); 066 TSourceToken endToken = gspObject.getEndToken(); 067 if (startToken != null) { 068 this.startPosition = new Pair3<Long, Long, String>(startToken.lineNo, startToken.columnNo, 069 ModelBindingManager.getGlobalHash()); 070 } 071 072 if (endToken != null) { 073 this.endPosition = new Pair3<Long, Long, String>(endToken.lineNo, 074 endToken.columnNo + SQLUtil.endTrim(endToken.getAstext()).length(), ModelBindingManager.getGlobalHash()); 075 } 076 077 this.schema = ModelBindingManager.getGlobalSchema(); 078 this.database = ModelBindingManager.getGlobalDatabase(); 079 this.queryHashId = SQLUtil.stringToMD5(gspObject.asCanonical().replaceAll("\r?\n", "\r\n")); 080 081 EDbVendor vendor = ModelBindingManager.getGlobalOption().getVendor(); 082 boolean supportCatalog = TSQLEnv.supportCatalog(vendor); 083 boolean supportSchema = TSQLEnv.supportSchema(vendor); 084 085 fillSchemaInfo(); 086 087 if (!supportCatalog) { 088 this.database = null; 089 } else if (this.database == null && !TSQLEnv.DEFAULT_DB_NAME.equals(getDefaultDatabase())) { 090 this.database = getDefaultDatabase(); 091 } 092 093 if (!supportSchema) { 094 this.schema = null; 095 } else if (this.schema == null && !TSQLEnv.DEFAULT_SCHEMA_NAME.equals(getDefaultSchema())) { 096 this.schema = getDefaultSchema(); 097 } 098 099 String procedureParent = getProcedureParentName(gspObject); 100 if (procedureParent != null) { 101 procedureName = procedureParent; 102 Procedure procedure = ModelBindingManager.get() 103 .getProcedureByName(DlineageUtil.getIdentifierNormalTableName( 104 DlineageUtil.getProcedureNameWithArgs(getParentProcedure(gspObject)))); 105 if (procedure == null) { 106 procedure = ModelBindingManager.get().getProcedureByName(DlineageUtil.getIdentifierNormalTableName( 107 DlineageUtil.getProcedureNameWithArgNum(getParentProcedure(gspObject)))); 108 } 109 if (procedure != null) { 110 this.procedureId = procedure.getId(); 111 } 112 } else { 113 procedureName = "batchQueries"; 114 } 115 116 TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek()); 117 String sqlComment = null; 118 try { 119 sqlComment = stmt.getCommentBeforeNode(); 120 } catch (Exception e) { 121 } 122 if (!SQLUtil.isEmpty(sqlComment) && sqlComment.indexOf("process") != -1) { 123 Properties properties = new Properties(); 124 try { 125 properties.load( 126 new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes())); 127 if (properties.containsKey("process_label")) { 128 this.customType = properties.getProperty("process_label"); 129 } else if (properties.containsKey("process")) { 130 this.customType = properties.getProperty("process"); 131 } 132 } catch (IOException e) { 133 logger.error("load sql comment properties failed.", e); 134 } 135 } 136 137 if (gspObject instanceof TCreateTableSqlStatement) { 138 this.type = "Create Table"; 139 } else if (gspObject instanceof TCreateViewSqlStatement) { 140 this.type = "Create View"; 141 } else if (gspObject instanceof TUpdateSqlStatement) { 142 this.type = "Update"; 143 } else if (gspObject instanceof TMergeSqlStatement) { 144 this.type = "Merge"; 145 } else if (gspObject instanceof TInsertSqlStatement) { 146 this.type = "Insert"; 147 } else if (gspObject instanceof TSelectSqlStatement 148 && ((TSelectSqlStatement) gspObject).getIntoClause() != null) { 149 this.type = "Select Into"; 150 } else if (gspObject instanceof TCreateTriggerStmt) { 151 this.type = "Create Trigger"; 152 } else if (gspObject instanceof TCreateFunctionStmt) { 153 this.type = "Create Function"; 154 } else if (gspObject instanceof TMssqlCreateFunction) { 155 this.type = "Create Function"; 156 } else if (gspObject instanceof TCreateProcedureStmt) { 157 this.type = "Create Procedure"; 158 } else if (gspObject instanceof THiveLoad) { 159 this.type = "Hive Load"; 160 } else if (gspObject instanceof TSnowflakeCopyIntoStmt) { 161 this.type = "Copy Into"; 162 } else if (gspObject instanceof TAlterTableStatement) { 163 this.type = "Alter Table"; 164 } else if (gspObject instanceof TRenameStmt) { 165 this.type = "Rename Table"; 166 } else if (gspObject instanceof TMssqlDeclare) { 167 this.type = "Mssql Declare"; 168 } else if (gspObject instanceof TCreateStageStmt) { 169 this.type = "Create Stage"; 170 } else if (gspObject instanceof TCreateSynonymStmt) { 171 this.type = "Create Synonym"; 172 } else if (gspObject instanceof TCreateExternalDataSourceStmt) { 173 this.type = "Create Datasource"; 174 } else if (gspObject instanceof TCreateStreamStmt) { 175 this.type = "Create Stream"; 176 } else if (gspObject instanceof TCreateDatabaseSqlStatement) { 177 this.type = "Create Database"; 178 } else if (gspObject instanceof TCreateSchemaSqlStatement) { 179 this.type = "Create Schema"; 180 } else if (gspObject instanceof TMssqlExecute && ((TMssqlExecute)gspObject).getModuleName().toString().equalsIgnoreCase("sp_rename")) { 181 this.type = "Rename Table"; 182 } else if (gspObject instanceof TCallStatement){ 183 this.type = "Function Call"; 184 } else { 185 this.type = gspObject.getClass().getSimpleName(); 186 } 187 188 if (this.server == null && !TSQLEnv.DEFAULT_SERVER_NAME.equals(getDefaultServer())) { 189 this.server = getDefaultServer(); 190 } 191 192 if (stmt.getStatements() != null) { 193 for (int i = 0; i < stmt.getStatements().size(); i++) { 194 TCustomSqlStatement subquery = stmt.getStatements().get(i); 195 if (subquery instanceof TSelectSqlStatement) { 196 TSelectSqlStatement select = (TSelectSqlStatement) subquery; 197 appendTransform(select, ESetOperatorType.none); 198 } 199 } 200 } 201 } 202 203 public Process(TFunctionCall gspObject) { 204 205 if (gspObject == null) { 206 throw new IllegalArgumentException("Process arguments can't be null."); 207 } 208 209 id = ++ModelBindingManager.get().TABLE_COLUMN_ID; 210 211 212 TSourceToken startToken = gspObject.getStartToken(); 213 TSourceToken endToken = gspObject.getEndToken(); 214 if (startToken != null) { 215 this.startPosition = new Pair3<Long, Long, String>(startToken.lineNo, startToken.columnNo, 216 ModelBindingManager.getGlobalHash()); 217 } 218 219 if (endToken != null) { 220 this.endPosition = new Pair3<Long, Long, String>(endToken.lineNo, 221 endToken.columnNo + SQLUtil.endTrim(endToken.getAstext()).length(), ModelBindingManager.getGlobalHash()); 222 } 223 224 TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek()); 225 this.gspObject = stmt; 226 227 this.schema = ModelBindingManager.getGlobalSchema(); 228 this.database = ModelBindingManager.getGlobalDatabase(); 229 230 if (stmt != null) { 231 this.queryHashId = SQLUtil.stringToMD5(stmt.asCanonical().replaceAll("\r?\n", "\r\n")); 232 } 233 234 EDbVendor vendor = ModelBindingManager.getGlobalOption().getVendor(); 235 boolean supportCatalog = TSQLEnv.supportCatalog(vendor); 236 boolean supportSchema = TSQLEnv.supportSchema(vendor); 237 238 fillSchemaInfo(); 239 240 if (!supportCatalog) { 241 this.database = null; 242 } else if (this.database == null && !TSQLEnv.DEFAULT_DB_NAME.equals(getDefaultDatabase())) { 243 this.database = getDefaultDatabase(); 244 } 245 246 if (!supportSchema) { 247 this.schema = null; 248 } else if (this.schema == null && !TSQLEnv.DEFAULT_SCHEMA_NAME.equals(getDefaultSchema())) { 249 this.schema = getDefaultSchema(); 250 } 251 252 String procedureParent = getProcedureParentName(stmt); 253 if (procedureParent != null) { 254 procedureName = procedureParent; 255 Procedure procedure = ModelBindingManager.get() 256 .getProcedureByName(DlineageUtil.getIdentifierNormalTableName( 257 DlineageUtil.getProcedureNameWithArgs(getParentProcedure(stmt)))); 258 if (procedure == null) { 259 procedure = ModelBindingManager.get().getProcedureByName(DlineageUtil.getIdentifierNormalTableName( 260 DlineageUtil.getProcedureNameWithArgNum(getParentProcedure(stmt)))); 261 } 262 if (procedure != null) { 263 this.procedureId = procedure.getId(); 264 } 265 } else { 266 procedureName = "batchQueries"; 267 } 268 269 String sqlComment = null; 270 try { 271 sqlComment = stmt.getCommentBeforeNode(); 272 } catch (Exception e) { 273 } 274 if (!SQLUtil.isEmpty(sqlComment) && sqlComment.indexOf("process") != -1) { 275 Properties properties = new Properties(); 276 try { 277 properties.load( 278 new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes())); 279 if (properties.containsKey("process_label")) { 280 this.customType = properties.getProperty("process_label"); 281 } else if (properties.containsKey("process")) { 282 this.customType = properties.getProperty("process"); 283 } 284 } catch (IOException e) { 285 logger.error("load sql comment properties failed.", e); 286 } 287 } 288 289 this.type = "Function Call"; 290 291 if (this.server == null && !TSQLEnv.DEFAULT_SERVER_NAME.equals(getDefaultServer())) { 292 this.server = getDefaultServer(); 293 } 294 } 295 296 private void appendTransform(TSelectSqlStatement select, ESetOperatorType operatorType) { 297 // Iterative traversal of the UNION tree to avoid StackOverflow with deeply nested unions. 298 Deque<Object[]> stack = new ArrayDeque<>(); 299 stack.push(new Object[]{select, operatorType}); 300 301 while (!stack.isEmpty()) { 302 Object[] entry = stack.pop(); 303 TSelectSqlStatement current = (TSelectSqlStatement) entry[0]; 304 ESetOperatorType currentOp = (ESetOperatorType) entry[1]; 305 306 if (currentOp != ESetOperatorType.none) { 307 Transform transform = new Transform(); 308 transform.setType(currentOp.name()); 309 transform.setCodeString(current.toString()); 310 transform.setStartToken(current.getStartToken()); 311 transform.setEndToken(current.getEndToken()); 312 transforms.add(transform); 313 } else { 314 if (current.getCteList() != null && current.getCteList().size() > 0) { 315 for (int i = 0; i < current.getCteList().size(); i++) { 316 TCTE cte = current.getCteList().getCTE(i); 317 Transform transform = new Transform(); 318 transform.setType(Transform.CTE); 319 transform.setCodeString(cte.toString()); 320 transform.setStartToken(cte.getStartToken()); 321 transform.setEndToken(cte.getEndToken()); 322 transforms.add(transform); 323 } 324 } 325 if (current.getSetOperatorType() != ESetOperatorType.none) { 326 // Push right first so left is processed first (stack is LIFO) 327 if (current.getRightStmt() != null) { 328 ESetOperatorType rightOp = current.getRightStmt().getSetOperatorType() != ESetOperatorType.none 329 ? ESetOperatorType.none : current.getSetOperatorType(); 330 stack.push(new Object[]{current.getRightStmt(), rightOp}); 331 } 332 if (current.getLeftStmt() != null) { 333 ESetOperatorType leftOp = current.getLeftStmt().getSetOperatorType() != ESetOperatorType.none 334 ? ESetOperatorType.none : current.getSetOperatorType(); 335 stack.push(new Object[]{current.getLeftStmt(), leftOp}); 336 } 337 } else { 338 if (current.getJoins() != null && current.getJoins().size() > 0) { 339 TJoinList joins = current.getJoins(); 340 TSourceToken startToken = joins.getStartToken(); 341 TSourceToken fromToken = SourceTokenSearcher.backforwardSearch(startToken, 10, "from"); 342 if (fromToken == null) { 343 if (startToken.getAstext().equalsIgnoreCase("from")) { 344 fromToken = startToken; 345 } else { 346 continue; 347 } 348 } 349 StringBuilder builder = new StringBuilder(); 350 if (joins.getEndToken().posinlist < fromToken.container.size() 351 && joins.getEndToken().posinlist >= fromToken.posinlist) { 352 for (int j = fromToken.posinlist; j <= joins.getEndToken().posinlist; j++) { 353 builder.append(fromToken.container.get(j)); 354 } 355 } else { 356 System.err.println("Handle statement transform error, statemet is:"); 357 System.err.println(current.toString()); 358 } 359 Transform transform = new Transform(); 360 transform.setType(Transform.FROM); 361 transform.setCodeString(builder.toString()); 362 transform.setStartToken(fromToken); 363 transform.setEndToken(joins.getEndToken()); 364 transforms.add(transform); 365 } 366 } 367 } 368 } 369 } 370 371 private TSelectSqlStatement getFirstSubquery(TSelectSqlStatement select) { 372 // Iterative left-chain walk to avoid StackOverflow with deeply nested unions. 373 TSelectSqlStatement current = select; 374 while (current.getSetOperatorType() != ESetOperatorType.none) { 375 current = current.getLeftStmt(); 376 } 377 return current; 378 } 379 380 private TSelectSqlStatement getLastSubquery(TSelectSqlStatement select) { 381 // Iterative right-chain walk to avoid StackOverflow with deeply nested unions. 382 TSelectSqlStatement current = select; 383 while (current.getSetOperatorType() != ESetOperatorType.none) { 384 current = current.getRightStmt(); 385 } 386 return current; 387 } 388 389 private void fillSchemaInfo() { 390 TCustomSqlStatement stmt = DlineageUtil.getTopStmt(ModelBindingManager.getGlobalStmtStack().peek()); 391 String sqlComment = null; 392 try { 393 sqlComment = stmt.getCommentBeforeNode(); 394 } catch (Exception e) { 395 } 396 if (!SQLUtil.isEmpty(sqlComment) && (sqlComment.indexOf("db") != -1 || sqlComment.indexOf("schema") != -1)) { 397 Properties properties = new Properties(); 398 try { 399 properties.load( 400 new ByteArrayInputStream(sqlComment.replace("--", "").trim().replace(",", "\n").getBytes())); 401 if (SQLUtil.isEmpty(this.server) && properties.containsKey("db-instance")) { 402 this.server = properties.getProperty("db-instance"); 403 } 404 if (SQLUtil.isEmpty(this.database) && properties.containsKey("db")) { 405 this.database = properties.getProperty("db"); 406 if (this.database.indexOf(".") != -1) { 407 this.database = SQLUtil.quoteDottedName(ModelBindingManager.getGlobalOption().getVendor(), ESQLDataObjectType.dotCatalog, this.database); 408 } 409 } 410 if (SQLUtil.isEmpty(this.schema) && properties.containsKey("schema")) { 411 this.schema = properties.getProperty("schema"); 412 if (this.schema.indexOf(".") != -1) { 413 this.schema = SQLUtil.quoteDottedName(ModelBindingManager.getGlobalOption().getVendor(), ESQLDataObjectType.dotSchema, this.schema); 414 } 415 } 416 } catch (IOException e) { 417 logger.error("load sql comment properties failed.", e); 418 } 419 } 420 } 421 422 protected String getDefaultServer() { 423 String defaultServer = null; 424 if (ModelBindingManager.getGlobalSQLEnv() != null) { 425 defaultServer = ModelBindingManager.getGlobalSQLEnv().getDefaultServerName(); 426 } 427 if (!SQLUtil.isEmpty(defaultServer)) 428 return defaultServer; 429 return TSQLEnv.DEFAULT_SERVER_NAME; 430 } 431 432 protected String getDefaultSchema() { 433 return ModelBindingManager.getEffectiveDefaultSchema(); 434 } 435 436 protected String getDefaultDatabase() { 437 return ModelBindingManager.getEffectiveDefaultDatabase(); 438 } 439 440 private String getProcedureParentName(TCustomSqlStatement stmt) { 441 if (stmt instanceof TStoredProcedureSqlStatement) { 442 if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) { 443 return ((TStoredProcedureSqlStatement) stmt).getStoredProcedureName().toString(); 444 } 445 } 446 447 if (stmt == null) { 448 return null; 449 } 450 451 if (stmt instanceof TTeradataCreateProcedure) { 452 if (((TTeradataCreateProcedure) stmt).getProcedureName() != null) { 453 return ((TTeradataCreateProcedure) stmt).getProcedureName().toString(); 454 } 455 } 456 457 if (stmt instanceof TStoredProcedureSqlStatement) { 458 if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) { 459 return ((TStoredProcedureSqlStatement) stmt).getStoredProcedureName().toString(); 460 } 461 462 if (stmt instanceof TCommonBlock) { 463 if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) { 464 stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName(); 465 } 466 else { 467 return getProcedureParentName(stmt.getParentStmt()); 468 } 469 } 470 471 return null; 472 } 473 474 stmt = stmt.getParentStmt(); 475 if (stmt == null) { 476 if(ModelBindingManager.getGlobalProcedure()!=null) { 477 return getProcedureParentName(ModelBindingManager.getGlobalProcedure()); 478 } 479 else { 480 return null; 481 } 482 } 483 484 485 if (stmt instanceof TCommonBlock) { 486 if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) { 487 stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName(); 488 } 489 } 490 491 return getProcedureParentName(stmt); 492 } 493 494 private TStoredProcedureSqlStatement getParentProcedure(TCustomSqlStatement stmt) { 495 if (stmt == null) { 496 return null; 497 } 498 499 if (stmt instanceof TStoredProcedureSqlStatement) { 500 if (((TStoredProcedureSqlStatement) stmt).getStoredProcedureName() != null) { 501 return ((TStoredProcedureSqlStatement) stmt); 502 } 503 504 if (stmt instanceof TCommonBlock) { 505 if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) { 506 return (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName(); 507 } 508 else { 509 return getParentProcedure(stmt.getParentStmt()); 510 } 511 } 512 513 return null; 514 } 515 516 stmt = stmt.getParentStmt(); 517 if (stmt == null) { 518 if(ModelBindingManager.getGlobalProcedure()!=null) { 519 return ModelBindingManager.getGlobalProcedure(); 520 } 521 else { 522 return null; 523 } 524 } 525 526 if (stmt instanceof TCommonBlock) { 527 if (((TCommonBlock) stmt).getBlockBody().getParentObjectName() instanceof TStoredProcedureSqlStatement) { 528 stmt = (TStoredProcedureSqlStatement) ((TCommonBlock) stmt).getBlockBody().getParentObjectName(); 529 } 530 } 531 532 return getParentProcedure(stmt); 533 } 534 535 public Pair3<Long, Long, String> getStartPosition() { 536 return startPosition; 537 } 538 539 public Pair3<Long, Long, String> getEndPosition() { 540 return endPosition; 541 } 542 543 public TCustomSqlStatement getGspObject() { 544 return gspObject; 545 } 546 547 public long getId() { 548 return id; 549 } 550 551 public String getSchema() { 552 return schema; 553 } 554 555 public String getDatabase() { 556 return database; 557 } 558 559 public String getName() { 560 return name; 561 } 562 563 public String getType() { 564 return type; 565 } 566 567 public void setName(String name) { 568 this.name = name; 569 } 570 571 public void setType(String type) { 572 this.type = type; 573 } 574 575 public String getProcedureName() { 576 return procedureName; 577 } 578 579 public void setProcedureName(String procedureName) { 580 this.procedureName = procedureName; 581 } 582 583 public String getQueryHashId() { 584 return queryHashId; 585 } 586 587 public void setQueryHashId(String queryHashId) { 588 this.queryHashId = queryHashId; 589 } 590 591 public void appendColumn(Object tableColumn) { 592 if (tableColumn != null && !targetColumns.contains(tableColumn)) { 593 targetColumns.add(tableColumn); 594 } 595 } 596 597 public List<Object> getTargetColumns() { 598 return targetColumns; 599 } 600 601 public String getCustomType() { 602 return customType; 603 } 604 605 public Long getProcedureId() { 606 return procedureId; 607 } 608 609 public String getServer() { 610 return server; 611 } 612 613 public List<Transform> getTransforms() { 614 return transforms; 615 } 616}