Scrivere il risultato della query SQL su file di Apache Flink

Sep 08 2020

Ho il seguente compito:

  1. Creare un lavoro con richiesta SQL alla tabella Hive;
  2. Esegui questo lavoro sul cluster Flink remoto;
  3. 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

1 JarkWu Sep 08 2020 at 09:11

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