Tulis hasil SQL Query ke file oleh Apache Flink
Saya memiliki tugas berikut:
- Buat pekerjaan dengan permintaan SQL ke tabel Hive;
- Jalankan tugas ini di kluster Flink jarak jauh;
- 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
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