您好,登錄后才能下訂單哦!
今天寫了一個(gè)稍微復(fù)雜的例子, 實(shí)現(xiàn)了類似mysql group_concat 功能,記錄一下
MapToString 參考bug 那篇博客
public static void main(String[] arg) throws Exception {
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
BatchTableEnvironment tableEnv = new BatchTableEnvironment(env, TableConfig.DEFAULT());
tableEnv.registerFunction("mapToString", new MapToString());
getProjectInfo(env,tableEnv);
getProject(env,tableEnv);
joinTableProjectWithInfo(tableEnv);
Table query = tableEnv.sqlQuery("select id, name, type from result_agg");
DataSet<Row> ds= tableEnv.toDataSet(query, Row.class);
ds.print();
ds.writeAsText("/home/test", WriteMode.OVERWRITE);
env.execute("multiple-table");
}
public static void getProjectInfo(ExecutionEnvironment env,BatchTableEnvironment tableEnv) {
TypeInformation[] fieldTypes = new TypeInformation[] { BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO };
String[] fieldNames = new String[] { "id", "type" };
RowTypeInfo rowTypeInfo = new RowTypeInfo(fieldTypes, fieldNames);
JDBCInputFormat jdbcInputFormat = JDBCInputFormat.buildJDBCInputFormat().setDrivername("com.mysql.jdbc.Driver")
.setDBUrl("jdbc:mysql://ip:3306/space?characterEncoding=utf8")
.setUsername("user").setPassword("pwd")
.setQuery("select project_fid, cast(project_info_type as CHAR) as type from project").setRowTypeInfo(rowTypeInfo).finish();
DataSource<Row> s = env.createInput(jdbcInputFormat);
tableEnv.registerDataSet("project_info", s);
aggProjectInfo(tableEnv,"project_info");
}
public static void aggProjectInfo(BatchTableEnvironment tableEnv, String tableName) {
Table tapiResult = tableEnv.scan(tableName);
tapiResult.printSchema();
Table query = tableEnv.sqlQuery("select id, mapToString(collect(type)) as type from project_info group by id");
tableEnv.registerTable(tableName+"_agg", query);
tapiResult = tableEnv.scan(tableName+"_agg");
tapiResult.printSchema();
}
public static void getProject(ExecutionEnvironment env,BatchTableEnvironment tableEnv) {
TypeInformation[] fieldTypes = new TypeInformation[] { BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO };
String[] fieldNames = new String[] { "pid", "name" };
RowTypeInfo rowTypeInfo = new RowTypeInfo(fieldTypes, fieldNames);
JDBCInputFormat jdbcInputFormat = JDBCInputFormat.buildJDBCInputFormat().setDrivername("com.mysql.jdbc.Driver")
.setDBUrl("jdbc:mysql://ip:3306/space?characterEncoding=utf8")
.setUsername("user").setPassword("pwd")
.setQuery("select fid, project_name from t_project").setRowTypeInfo(rowTypeInfo).finish();
DataSource<Row> s = env.createInput(jdbcInputFormat);
tableEnv.registerDataSet("project", s);
}
public static void joinTableProjectWithInfo(BatchTableEnvironment tableEnv) {
Table result =tableEnv.sqlQuery("select a.pid as id , a.name , b.type from project a inner join project_info_agg b on a.pid=b.id");
tableEnv.registerTable("result_agg", result);
result.printSchema();
}
免責(zé)聲明:本站發(fā)布的內(nèi)容(圖片、視頻和文字)以原創(chuàng)、轉(zhuǎn)載和分享為主,文章觀點(diǎn)不代表本網(wǎng)站立場,如果涉及侵權(quán)請聯(lián)系站長郵箱:is@yisu.com進(jìn)行舉報(bào),并提供相關(guān)證據(jù),一經(jīng)查實(shí),將立刻刪除涉嫌侵權(quán)內(nèi)容。