001package gudusoft.gsqlparser.sqlcmds; 002 003import gudusoft.gsqlparser.*; 004import gudusoft.gsqlparser.stmt.*; 005import gudusoft.gsqlparser.stmt.sparksql.*; 006import gudusoft.gsqlparser.stmt.mysql.TLoadDataStmt; 007import gudusoft.gsqlparser.stmt.flink.TFlinkExplainStmt; 008import gudusoft.gsqlparser.stmt.flink.TFlinkCreateCatalogStmt; 009import gudusoft.gsqlparser.stmt.flink.TFlinkDropCatalogStmt; 010import gudusoft.gsqlparser.stmt.flink.TFlinkAlterCatalogStmt; 011import gudusoft.gsqlparser.stmt.flink.TFlinkExecuteStatementSetStmt; 012 013/** 014 * Apache Flink SQL command resolver. 015 * Based on SparkSQL implementation with Flink-specific additions. 016 * 017 * @since 3.2.0.0 018 */ 019public class TSqlCmdsFlink extends AbstractSqlCmds { 020 021 public TSqlCmdsFlink() { 022 super(EDbVendor.dbvflink); 023 } 024 025 @Override 026 protected void initializeCommands() { 027 // Commands ordered with longer patterns before shorter ones 028 029 // ALTER commands 030 addCmd(TBaseType.rrw_alter, "temporary", "system", "function", ESqlStatementType.sstalterfunction); 031 addCmd(TBaseType.rrw_alter, "temporary", "function", ESqlStatementType.sstalterfunction); 032 addCmd(TBaseType.rrw_alter, "catalog", ESqlStatementType.sstaltercatalog); 033 addCmd(TBaseType.rrw_alter, "database", ESqlStatementType.sstalterdatabase); 034 addCmd(TBaseType.rrw_alter, "table", ESqlStatementType.sstaltertable); 035 addCmd(TBaseType.rrw_alter, "view", ESqlStatementType.sstalterview); 036 addCmd(TBaseType.rrw_alter, "function", ESqlStatementType.sstalterfunction); 037 038 // ANALYZE commands 039 addCmd(TBaseType.rrw_analyze, "table", ESqlStatementType.sstanalyzeTable); 040 041 // CREATE commands (longer patterns first) 042 addCmd(TBaseType.rrw_create, "or", "replace", "temporary", "view", ESqlStatementType.sstcreateview); 043 addCmd(TBaseType.rrw_create, "or", "replace", "view", ESqlStatementType.sstcreateview); 044 addCmd(TBaseType.rrw_create, "or", "replace", "table", ESqlStatementType.sstcreatetable); 045 addCmd(TBaseType.rrw_create, "temporary", "system", "function", ESqlStatementType.sstcreatefunction); 046 addCmd(TBaseType.rrw_create, "temporary", "function", ESqlStatementType.sstcreatefunction); 047 addCmd(TBaseType.rrw_create, "temporary", "view", ESqlStatementType.sstcreateview); 048 addCmd(TBaseType.rrw_create, "temporary", "table", ESqlStatementType.sstcreatetable); 049 addCmd(TBaseType.rrw_create, "catalog", ESqlStatementType.sstcreatecatalog); 050 addCmd(TBaseType.rrw_create, "database", ESqlStatementType.sstcreatedatabase); 051 addCmd(TBaseType.rrw_create, "function", ESqlStatementType.sstcreatefunction); 052 addCmd(TBaseType.rrw_create, "table", ESqlStatementType.sstcreatetable); 053 addCmd(TBaseType.rrw_create, "view", ESqlStatementType.sstcreateview); 054 055 // DELETE commands 056 addCmd(TBaseType.rrw_delete, ESqlStatementType.sstdelete); 057 058 // DESCRIBE commands (DESC and DESCRIBE variants) 059 addCmd(TBaseType.rrw_spark_desc, "database", ESqlStatementType.sstdescribeDatabase); 060 addCmd(TBaseType.rrw_spark_desc, "function", ESqlStatementType.sstdescribeFunction); 061 addCmd(TBaseType.rrw_spark_desc, "table", ESqlStatementType.sstdescribeTable); 062 addCmd(TBaseType.rrw_spark_desc, ESqlStatementType.sstdescribeTable); 063 addCmd(TBaseType.rrw_describe, "database", ESqlStatementType.sstdescribeDatabase); 064 addCmd(TBaseType.rrw_describe, "function", ESqlStatementType.sstdescribeFunction); 065 addCmd(TBaseType.rrw_describe, "table", ESqlStatementType.sstdescribeTable); 066 addCmd(TBaseType.rrw_describe, ESqlStatementType.sstdescribe); 067 068 // DROP commands 069 addCmd(TBaseType.rrw_drop, "temporary", "system", "function", ESqlStatementType.sstdropfunction); 070 addCmd(TBaseType.rrw_drop, "temporary", "function", ESqlStatementType.sstdropfunction); 071 addCmd(TBaseType.rrw_drop, "temporary", "view", ESqlStatementType.sstdropview); 072 addCmd(TBaseType.rrw_drop, "temporary", "table", ESqlStatementType.sstdroptable); 073 addCmd(TBaseType.rrw_drop, "catalog", ESqlStatementType.sstdropcatalog); 074 addCmd(TBaseType.rrw_drop, "database", ESqlStatementType.sstdropdatabase); 075 addCmd(TBaseType.rrw_drop, "function", ESqlStatementType.sstdropfunction); 076 addCmd(TBaseType.rrw_drop, "table", ESqlStatementType.sstdroptable); 077 addCmd(TBaseType.rrw_drop, "view", ESqlStatementType.sstdropview); 078 079 // EXPLAIN commands 080 addCmd(TBaseType.rrw_explain, ESqlStatementType.sstExplain); 081 082 // INSERT commands 083 addCmd(TBaseType.rrw_insert, "overwrite", ESqlStatementType.sstinsert); 084 addCmd(TBaseType.rrw_insert, "into", ESqlStatementType.sstinsert); 085 addCmd(TBaseType.rrw_insert, ESqlStatementType.sstinsert); 086 087 // REPLACE commands (Flink RTAS) 088 addCmd(TBaseType.rrw_replace, "table", ESqlStatementType.sstcreatetable); 089 090 // RESET commands 091 addCmd(TBaseType.rrw_reset, ESqlStatementType.sstReset); 092 093 // SELECT commands 094 addCmd(TBaseType.rrw_select, ESqlStatementType.sstselect); 095 096 // SET commands 097 addCmd(TBaseType.rrw_set, ESqlStatementType.sstset); 098 099 // SHOW commands 100 addCmd(TBaseType.rrw_show, "create", "table", ESqlStatementType.sstShowCreateTable); 101 addCmd(TBaseType.rrw_show, "create", "view", ESqlStatementType.sstShow); 102 addCmd(TBaseType.rrw_show, "create", "catalog", ESqlStatementType.sstShow); 103 addCmd(TBaseType.rrw_show, "user", "functions", ESqlStatementType.sstShowUserFunctions); 104 addCmd(TBaseType.rrw_show, "catalogs", ESqlStatementType.sstShow); 105 addCmd(TBaseType.rrw_show, "current", ESqlStatementType.sstShow); 106 addCmd(TBaseType.rrw_show, "modules", ESqlStatementType.sstShow); 107 addCmd(TBaseType.rrw_show, "full", "modules", ESqlStatementType.sstShow); 108 addCmd(TBaseType.rrw_show, "jars", ESqlStatementType.sstShow); 109 addCmd(TBaseType.rrw_show, "jobs", ESqlStatementType.sstShow); 110 addCmd(TBaseType.rrw_show, "columns", ESqlStatementType.sstShowColumns); 111 addCmd(TBaseType.rrw_show, "databases", ESqlStatementType.sstShowDatabases); 112 addCmd(TBaseType.rrw_show, "functions", ESqlStatementType.sstShowFunctions); 113 addCmd(TBaseType.rrw_show, "partitions", ESqlStatementType.sstShowPartitions); 114 addCmd(TBaseType.rrw_show, "tables", ESqlStatementType.sstShowTables); 115 addCmd(TBaseType.rrw_show, "views", ESqlStatementType.sstShowViews); 116 117 // EXECUTE commands (Flink EXECUTE STATEMENT SET / EXECUTE PLAN) 118 addCmd(TBaseType.rrw_execute, ESqlStatementType.sstExecute); 119 120 // TRUNCATE commands 121 addCmd(TBaseType.rrw_truncate, "table", ESqlStatementType.ssttruncatetable); 122 123 // UPDATE commands 124 addCmd(TBaseType.rrw_update, ESqlStatementType.sstupdate); 125 126 // USE commands 127 addCmd(TBaseType.rrw_use, ESqlStatementType.sstUse); 128 } 129 130 @Override 131 protected String getToken1Str(int token1) { 132 // Flink vendor-specific tokens (inherited from SparkSQL) 133 switch (token1) { 134 case TBaseType.rrw_spark_desc: 135 return "desc"; 136 default: 137 return null; 138 } 139 } 140 141 @Override 142 public TCustomSqlStatement issql(TSourceToken token, EFindSqlStateType state, TCustomSqlStatement currentStatement) { 143 144 TCustomSqlStatement ret = null; 145 int k; 146 boolean lcisnewsql; 147 TSourceToken lcpprevsolidtoken, lcnextsolidtoken; 148 149 gnewsqlstatementtype = ESqlStatementType.sstinvalid; 150 151 if ((token.tokencode == TBaseType.cmtdoublehyphen) 152 || (token.tokencode == TBaseType.cmtslashstar) 153 || (token.tokencode == TBaseType.lexspace) 154 || (token.tokencode == TBaseType.lexnewline) 155 || (token.tokentype == ETokenType.ttsemicolon)) { 156 return ret; 157 } 158 159 int lcpos = token.posinlist; 160 TSourceTokenList lcsourcetokenlist = token.container; 161 TCustomSqlStatement lccurrentsqlstatement = currentStatement; 162 163 // subquery after semicolon || at first line 164 if ((state == EFindSqlStateType.stnormal) && (token.tokentype == ETokenType.ttleftparenthesis)) { 165 k = lcsourcetokenlist.solidtokenafterpos(lcpos, TBaseType.rrw_select, 1, "("); 166 if (k > 0) { 167 ret = new TSelectSqlStatement(this.vendor); 168 } 169 return ret; 170 } 171 172 // cte 173 if ((state == EFindSqlStateType.stnormal) && (token.tokencode == TBaseType.rrw_with)) { 174 ret = findcte(token); 175 if ((ret != null)) return ret; 176 } 177 178 gnewsqlstatementtype = getStatementTypeForToken(token); 179 180 // Flink EXECUTE STATEMENT SET trailing END marker: a top-level statement 181 // that starts with END terminates a statement set. 182 if ((gnewsqlstatementtype == ESqlStatementType.sstinvalid) 183 && (token.tokencode == TBaseType.rrw_end) 184 && (state == EFindSqlStateType.stnormal)) { 185 gnewsqlstatementtype = ESqlStatementType.sstExecute; 186 return new TFlinkExecuteStatementSetStmt(this.vendor); 187 } 188 189 TSourceToken lcprevsolidtoken = lcsourcetokenlist.solidtokenbefore(lcpos); 190 191 if ((gnewsqlstatementtype == ESqlStatementType.sstinvalid) && (token.tokencode == TBaseType.rrw_create)) { 192 TSourceToken viewToken = token.container.searchToken(TBaseType.rrw_view, "", token, 15); 193 if (viewToken != null) { 194 gnewsqlstatementtype = ESqlStatementType.sstcreateview; 195 } 196 } 197 198 switch (gnewsqlstatementtype) { 199 case sstinvalid: { 200 ret = null; 201 break; 202 } 203 case sstselect: { 204 lcisnewsql = true; 205 206 if (state != EFindSqlStateType.stnormal) { 207 if (TBaseType.assigned(lcprevsolidtoken)) { 208 if (lcprevsolidtoken.tokentype == ETokenType.ttleftparenthesis) 209 lcisnewsql = false; // subquery 210 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_union) 211 lcisnewsql = false; 212 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_intersect) 213 lcisnewsql = false; 214 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_minus) 215 lcisnewsql = false; 216 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_except) 217 lcisnewsql = false; 218 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_return) 219 lcisnewsql = false; 220 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_as) { 221 if (lccurrentsqlstatement.sqlstatementtype == ESqlStatementType.sstcreatetable) 222 lcisnewsql = false; 223 if (lccurrentsqlstatement.sqlstatementtype == ESqlStatementType.sstcreateview) 224 lcisnewsql = false; 225 } 226 227 if (lcisnewsql && (lcprevsolidtoken.tokencode == TBaseType.rrw_all)) { 228 lcpprevsolidtoken = lcsourcetokenlist.solidtokenbefore(lcprevsolidtoken.posinlist); 229 if (TBaseType.assigned(lcpprevsolidtoken)) { 230 if (lcpprevsolidtoken.tokencode == TBaseType.rrw_union) 231 lcisnewsql = false; 232 } 233 } 234 } 235 236 if (TBaseType.assigned(lccurrentsqlstatement)) { 237 if (lccurrentsqlstatement.sqlstatementtype == ESqlStatementType.sstinsert) 238 lcisnewsql = false; 239 } 240 } 241 242 if (lcisnewsql) 243 ret = new TSelectSqlStatement(this.vendor); 244 break; 245 } 246 case sstinsert: { 247 lcisnewsql = true; 248 if (state != EFindSqlStateType.stnormal) { 249 if (TBaseType.assigned(lccurrentsqlstatement)) { 250 // No special handling needed 251 } 252 } 253 254 if (lcisnewsql) 255 ret = new TInsertSqlStatement(this.vendor); 256 ret.sqlstatementtype = gnewsqlstatementtype; 257 break; 258 } 259 case sstupdate: { 260 lcisnewsql = true; 261 if (state != EFindSqlStateType.stnormal) { 262 lcprevsolidtoken = lcsourcetokenlist.solidtokenbefore(lcpos); 263 if (TBaseType.assigned(lcprevsolidtoken)) { 264 if (lcprevsolidtoken.tokencode == TBaseType.rrw_on) 265 lcisnewsql = false; 266 else if (lcprevsolidtoken.tokencode == TBaseType.rrw_for) 267 lcisnewsql = false; 268 } 269 270 lcnextsolidtoken = lcsourcetokenlist.nextsolidtoken(lcpos, 1, false); 271 if (TBaseType.assigned(lcnextsolidtoken)) { 272 if (lcnextsolidtoken.tokentype == ETokenType.ttleftparenthesis) { 273 k = lcsourcetokenlist.solidtokenafterpos(lcnextsolidtoken.posinlist, TBaseType.rrw_select, 1, "("); 274 if (k == 0) lcisnewsql = false; 275 } 276 } 277 278 if (TBaseType.assigned(lccurrentsqlstatement)) { 279 // No special handling needed 280 } 281 } 282 283 if (lcisnewsql) { 284 ret = new TUpdateSqlStatement(this.vendor); 285 ret.dummytag = 1; 286 } 287 break; 288 } 289 case sstdelete: { 290 lcisnewsql = true; 291 292 if (state != EFindSqlStateType.stnormal) { 293 lcprevsolidtoken = lcsourcetokenlist.solidtokenbefore(lcpos); 294 if (TBaseType.assigned(lcprevsolidtoken)) { 295 if (lcprevsolidtoken.tokencode == TBaseType.rrw_on) 296 lcisnewsql = false; 297 } 298 299 if (TBaseType.assigned(lccurrentsqlstatement)) { 300 // No special handling needed 301 } 302 } 303 304 if (lcisnewsql) 305 ret = new TDeleteSqlStatement(this.vendor); 306 break; 307 } 308 case sstcreatetable: { 309 ret = new TCreateTableSqlStatement(this.vendor); 310 break; 311 } 312 case sstcreateview: { 313 ret = new TCreateViewSqlStatement(this.vendor); 314 break; 315 } 316 case sstcreatedatabase: { 317 ret = new TCreateDatabaseSqlStatement(this.vendor); 318 break; 319 } 320 case sstcreatecatalog: { 321 ret = new TFlinkCreateCatalogStmt(this.vendor); 322 break; 323 } 324 case sstdroptable: { 325 ret = new TDropTableSqlStatement(this.vendor); 326 break; 327 } 328 case sstdropview: { 329 ret = new TDropViewSqlStatement(this.vendor); 330 break; 331 } 332 case sstdropdatabase: { 333 ret = new TDropDatabaseStmt(this.vendor); 334 break; 335 } 336 case sstdropcatalog: { 337 ret = new TFlinkDropCatalogStmt(this.vendor); 338 break; 339 } 340 case sstaltertable: { 341 ret = new TAlterTableStatement(this.vendor); 342 break; 343 } 344 case sstalterview: { 345 ret = new TAlterViewStatement(this.vendor); 346 break; 347 } 348 case sstalterdatabase: { 349 ret = new TAlterDatabaseStmt(this.vendor); 350 break; 351 } 352 case sstaltercatalog: { 353 ret = new TFlinkAlterCatalogStmt(this.vendor); 354 break; 355 } 356 case sstset: 357 case sstReset: { 358 lcisnewsql = true; 359 if (state != EFindSqlStateType.stnormal) { 360 if (TBaseType.assigned(lccurrentsqlstatement)) { 361 lcisnewsql = false; 362 } 363 } 364 365 if (lcisnewsql) { 366 ret = new TSetStmt(this.vendor); 367 } 368 break; 369 } 370 case sstcreatefunction: { 371 ret = new TCreateFunctionStmt(this.vendor); 372 break; 373 } 374 case sstdropfunction: { 375 ret = new TDropFunctionStmt(this.vendor); 376 break; 377 } 378 case sstalterfunction: { 379 ret = new TAlterFunctionStmt(this.vendor); 380 break; 381 } 382 case sstTruncate: 383 case ssttruncatetable: { 384 ret = new TTruncateStatement(this.vendor); 385 break; 386 } 387 case sstdescribe: { 388 ret = new TDescribeStmt(this.vendor); 389 break; 390 } 391 case sstdescribeDatabase: 392 case sstdescribeTable: 393 case sstdescribeFunction: { 394 ret = new TDescribeStmt(this.vendor); 395 ret.sqlstatementtype = gnewsqlstatementtype; 396 break; 397 } 398 case sstExplain: { 399 ret = new TFlinkExplainStmt(this.vendor); 400 break; 401 } 402 case sstUse: { 403 ret = new TUseDatabase(this.vendor); 404 break; 405 } 406 case sstShow: 407 case sstShowColumns: 408 case sstShowCreateTable: 409 case sstShowDatabases: 410 case sstShowFunctions: 411 case sstShowUserFunctions: 412 case sstShowPartitions: 413 case sstShowTables: 414 case sstShowViews: { 415 ret = new TShowStmt(this.vendor); 416 ret.sqlstatementtype = gnewsqlstatementtype; 417 break; 418 } 419 case sstExecute: { 420 ret = new TFlinkExecuteStatementSetStmt(this.vendor); 421 break; 422 } 423 case sstanalyzeTable: { 424 ret = new TAnalyzeStmt(this.vendor); 425 break; 426 } 427 default: { 428 ret = new TUnknownSqlStatement(this.vendor); 429 ret.sqlstatementtype = gnewsqlstatementtype; 430 break; 431 } 432 } 433 434 return ret; 435 } 436}