Table In Flink
之前没有使用过 Table 相关的 API,在这稍微做下记录。
之前没有使用过 Table 相关的 API,在这稍微做下记录。
Structure of Program
一个使用 Table API 的程序结构大致如下所示:
// create a TableEnvironment for specific planner batch or streamingTableEnvironment tableEnv = ...; // see "Create a TableEnvironment" section
// register a TabletableEnv.registerTable("table1", ...) // ortableEnv.registerTableSource("table2", ...); // or// register an output TabletableEnv.registerTableSink("outputTable", ...);
// create a Table from a Table API queryTable tapiResult = tableEnv.scan("table1").select(...);// create a Table from a SQL queryTable sqlResult = tableEnv.sqlQuery("SELECT ... FROM table2 ... ");
// emit a Table API result Table to a TableSink, same for SQL resulttapiResult.insertInto("outputTable");
// executetableEnv.execute("java_job");TableEnvironment 分为 Batch 和 Stream 两种,分别在 bridge 模块中可以实现与 DataSet / DataStream 之间的转换。
注册 Table 的方式有三种:
- registerTable // register with Table
- registerTableSource // register with TableSource
- registerTableSink // register with TableSink
注册后便会在 Catalog 下注册对应的表,目前支持 Hive 和 InMemory 两种模式。