เขียนผลลัพธ์ของ SQL Query ไปยังไฟล์โดย Apache Flink
ฉันมีภารกิจต่อไปนี้:
- สร้างงานด้วยการร้องขอ SQL ไปยังตาราง Hive;
- รันงานนี้บนคลัสเตอร์ Flink ระยะไกล
- รวบรวมผลลัพธ์ของงานนี้ในไฟล์ (ควรใช้ HDFS)
บันทึก
เนื่องจากจำเป็นต้องรันงานนี้บนคลัสเตอร์ Flink ระยะไกลฉันจึงไม่สามารถใช้TableEnvironment ได้ด้วยวิธีง่ายๆ ปัญหานี้กล่าวถึงในตั๋วนี้:https://issues.apache.org/jira/browse/FLINK-18095. สำหรับโซลูชันปัจจุบันฉันใช้ adivce จากhttp://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/Table-Environment-for-Remote-Execution-td35691.html.
รหัส
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");
คำถาม
ในขั้นตอนนี้ฉันสามารถเรียกใช้วิธีtable.execute ()และหลังจากได้รับCloseableIteratorโดยวิธีการcollect () แต่ในกรณีของฉันฉันได้รับจำนวนแถวจำนวนมากอันเป็นผลมาจากคำขอของฉันและจะเป็นการดีที่จะรวบรวมเป็นไฟล์ (ORC ใน HDFS)
ฉันจะบรรลุเป้าหมายได้อย่างไร?
คำตอบ
Table.execute().collect()ส่งคืนผลลัพธ์ของมุมมองไปยังฝั่งไคลเอ็นต์ของคุณเพื่อวัตถุประสงค์ในการโต้ตอบ ในกรณีของคุณคุณสามารถใช้ตัวเชื่อมต่อระบบไฟล์และใช้INSERT INTOสำหรับเขียนมุมมองไปยังไฟล์ ตัวอย่างเช่น:
// 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");
ดูเพิ่มเติมเกี่ยวกับตัวเชื่อมต่อระบบไฟล์: https://ci.apache.org/projects/flink/flink-docs-release-1.11/dev/table/connectors/filesystem.html