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}