Tulis hasil SQL Query ke file oleh Apache Flink

Sep 08 2020

Saya memiliki tugas berikut:

  1. Buat pekerjaan dengan permintaan SQL ke tabel Hive;
  2. Jalankan tugas ini di kluster Flink jarak jauh;
  3. Kumpulkan hasil pekerjaan ini dalam file (HDFS lebih disukai).

Catatan

Karena itu perlu untuk menjalankan pekerjaan ini pada klaster Flink jarak jauh, saya tidak dapat menggunakan TableEnvironment dengan cara yang sederhana. Masalah ini disebutkan di tiket ini:https://issues.apache.org/jira/browse/FLINK-18095. Untuk solusi saat ini saya menggunakan adivce fromhttp://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Table-Environment-for-Remote-Execution-td35691.html.

Kode

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");

Pertanyaan

Pada langkah ini saya bisa memanggil () table.execute metode dan setelah mendapatkan CloseableIterator oleh mengumpulkan () metode. Tetapi dalam kasus saya, saya bisa mendapatkan jumlah baris yang besar sebagai hasil dari permintaan saya dan akan sempurna untuk mengumpulkannya ke dalam file (ORC dalam HDFS).

Bagaimana saya bisa mencapai tujuan saya?

Jawaban

1 JarkWu Sep 08 2020 at 09:11

Table.execute().collect()mengembalikan hasil tampilan ke sisi klien Anda untuk tujuan interaktif. Dalam kasus Anda, Anda dapat menggunakan konektor sistem file dan digunakan INSERT INTOuntuk menulis tampilan ke file. Sebagai contoh:

// 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");

Lihat lebih lanjut tentang konektor sistem file: https://ci.apache.org/projects/flink/flink-docs-release-1.11/dev/table/connectors/filesystem.html