Scrivere il risultato della query SQL su file di Apache Flink
Ho il seguente compito:
- Creare un lavoro con richiesta SQL alla tabella Hive;
- Esegui questo lavoro sul cluster Flink remoto;
- Raccogli il risultato di questo lavoro in un file (è preferibile HDFS).
Nota
Poiché è necessario eseguire questo lavoro sul cluster Flink remoto, non posso utilizzare TableEnvironment in modo semplice. Questo problema è menzionato in questo ticket:https://issues.apache.org/jira/browse/FLINK-18095. Per la soluzione corrente che uso adivce dahttp://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Table-Environment-for-Remote-Execution-td35691.html.
Codice
EnvironmentSettings batchSettings = EnvironmentSettings.newInstance().useBlinkPlanner().inBatchMode().build();
// create remote env
StreamExecutionEnvironment streamExecutionEnvironment = StreamExecutionEnvironment.createRemoteEnvironment("localhost", 8081, "/path/to/my/jar");
// create StreamTableEnvironment
TableConfig tableConfig = new TableConfig();
ClassLoader classLoader = Thread.currentThread().getContextClassLoader();
CatalogManager catalogManager = CatalogManager.newBuilder()
.classLoader(classLoader)
.config(tableConfig.getConfiguration())
.defaultCatalog(
batchSettings.getBuiltInCatalogName(),
new GenericInMemoryCatalog(
batchSettings.getBuiltInCatalogName(),
batchSettings.getBuiltInDatabaseName()))
.executionConfig(
streamExecutionEnvironment.getConfig())
.build();
ModuleManager moduleManager = new ModuleManager();
BatchExecutor batchExecutor = new BatchExecutor(streamExecutionEnvironment);
FunctionCatalog functionCatalog = new FunctionCatalog(tableConfig, catalogManager, moduleManager);
StreamTableEnvironmentImpl tableEnv = new StreamTableEnvironmentImpl(
catalogManager,
moduleManager,
functionCatalog,
tableConfig,
streamExecutionEnvironment,
new BatchPlanner(batchExecutor, tableConfig, functionCatalog, catalogManager),
batchExecutor,
false);
// configure HiveCatalog
String name = "myhive";
String defaultDatabase = "default";
String hiveConfDir = "/path/to/hive/conf"; // a local path
HiveCatalog hive = new HiveCatalog(name, defaultDatabase, hiveConfDir);
tableEnv.registerCatalog("myhive", hive);
tableEnv.useCatalog("myhive");
// request to Hive
Table table = tableEnv.sqlQuery("select * from myhive.`default`.test");
Domanda
In questo passaggio posso chiamare il metodo table.execute () e successivamente ottenere CloseableIterator con il metodo collect () . Ma nel mio caso posso ottenere un gran numero di righe come risultato della mia richiesta e sarà perfetto per raccoglierlo in file (ORC in HDFS).
Come posso raggiungere il mio obiettivo?
Risposte
Table.execute().collect()restituisce il risultato della visualizzazione al tuo lato client per scopi interattivi. Nel tuo caso, puoi usare il connettore del filesystem e usarlo INSERT INTOper scrivere la vista sul file. Per esempio:
// create a filesystem table
tableEnvironment.executeSql("CREATE TABLE MyUserTable (\n" +
" column_name1 INT,\n" +
" column_name2 STRING,\n" +
" ..." +
" \n" +
") WITH (\n" +
" 'connector' = 'filesystem',\n" +
" 'path' = 'hdfs://path/to/your/file',\n" +
" 'format' = 'orc' \n" +
")");
// submit the job
tableEnvironment.executeSql("insert into MyUserTable select * from myhive.`default`.test");
Ulteriori informazioni sul connettore del file system: https://ci.apache.org/projects/flink/flink-docs-release-1.11/dev/table/connectors/filesystem.html