Commit 8793b9c8 authored by 许育's avatar 许育

add: init

parents
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>org.example</groupId>
<artifactId>qm</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<java.version>1.8</java.version>
<org.mybatis.generator.version>3.4.6</org.mybatis.generator.version>
<skipTests>true</skipTests>
</properties>
<dependencies>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.22</version>
</dependency>
<dependency>
<groupId>org.apache.httpcomponents</groupId>
<artifactId>httpclient</artifactId>
<version>4.5.5</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.58</version>
</dependency>
<dependency>
<groupId>cglib</groupId>
<artifactId>cglib</artifactId>
<version>3.2.4</version>
</dependency>
<dependency>
<groupId>joda-time</groupId>
<artifactId>joda-time</artifactId>
<version>2.9.9</version>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.9</version>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.46</version>
<scope>compile</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>${java.version}</source>
<target>${java.version}</target>
<encoding>${project.build.sourceEncoding}</encoding>
</configuration>
</plugin>
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
<archive>
<manifest>
<mainClass>com.qm.doris.Load2Doris</mainClass>
</manifest>
<manifest>
<mainClass>com.qm.dorisv2.Load2DorisV2</mainClass>
</manifest>
<manifest>
<mainClass>com.qm.mysql.MysqlClient</mainClass>
</manifest>
</archive>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
\ No newline at end of file
This diff is collapsed.
package com.qm.doris;
import com.qm.model.BrokerLoadResult;
import com.qm.model.DescTable;
import com.qm.model.ShowCreateTable;
import com.qm.model.ShowPartitions;
import com.qm.util.MysqlJdbcUtil;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.Statement;
public class Load2DorisJdbc {
public static void executeVoidSql(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
stmt.executeUpdate(sql);
res = stmt.executeQuery("select 1");
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
}
public static BrokerLoadResult getStatusResult(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
return new BrokerLoadResult(String.valueOf(res.getString(1)),
String.valueOf(res.getString(2)),
String.valueOf(res.getString(3)),
String.valueOf(res.getString(4)),
String.valueOf(res.getString(5)),
String.valueOf(res.getString(6)),
String.valueOf(res.getString(7)),
String.valueOf(res.getString(8)),
String.valueOf(res.getString(9)),
String.valueOf(res.getString(10)),
String.valueOf(res.getString(11)),
String.valueOf(res.getString(12)),
String.valueOf(res.getString(13)),
String.valueOf(res.getString(14)),
String.valueOf(res.getString(15))
);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return null;
}
public static ShowPartitions getShowPartitionResult(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
return new ShowPartitions(String.valueOf(res.getString(1)),
String.valueOf(res.getString(2)),
String.valueOf(res.getString(3)),
String.valueOf(res.getString(4)),
String.valueOf(res.getString(5)),
String.valueOf(res.getString(6)),
String.valueOf(res.getString(7)),
String.valueOf(res.getString(8)),
String.valueOf(res.getString(9)),
String.valueOf(res.getString(10)),
String.valueOf(res.getString(11)),
String.valueOf(res.getString(12)),
String.valueOf(res.getString(13)),
String.valueOf(res.getString(14)),
String.valueOf(res.getString(15)),
String.valueOf(res.getString(16))
);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return null;
}
public static DescTable getDescResult(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
return new DescTable(String.valueOf(res.getString(1)),
String.valueOf(res.getString(2)),
String.valueOf(res.getString(3)),
String.valueOf(res.getString(4)),
String.valueOf(res.getString(5)),
String.valueOf(res.getString(6))
);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return null;
}
public static ShowCreateTable dorisPartitionTable(String sql, String url, String user, String password) {
ShowCreateTable showCreateTable = new ShowCreateTable();
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
String schmae = String.valueOf(res.getString(2));
if (schmae.contains("PARTITION BY RANGE(")) {
String[] partitionByRanges = schmae.split("PARTITION BY RANGE\\(");
showCreateTable.setPartition(partitionByRanges[1].split("\\)")[0]);
}
if (schmae.contains("\"dynamic_partition.enable\" = \"true\"")) {
showCreateTable.setDynamicPartition(true);
}
return showCreateTable;
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return showCreateTable;
}
}
This diff is collapsed.
package com.qm.dorisv2;
import com.qm.util.MysqlJdbcUtil;
import com.qm.util.SqlUtil;
import java.sql.Connection;
import java.sql.Statement;
import java.util.List;
public class Load2DorisJdbcV2 {
public static void executeVoidSql(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
List<String> sqls = SqlUtil.splitQueries(SqlUtil.sqlTrim(sql));
for (String sqll : sqls) {
if (stmt.execute(sqll)) {
// ResultSet resultSet = stmt.getResultSet();
// System.out.println("查询:" + sqll + "\n" + SqlUtil.resultSet2Array(resultSet));
} else {
// System.out.println("非查询:" + sqll + "\n" + stmt.getUpdateCount());
}
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
}
}
This diff is collapsed.
package com.qm.dorisv2;
import com.qm.model.Partition;
import com.qm.util.EmptyUtils;
import org.apache.commons.lang3.StringUtils;
import java.util.Arrays;
import java.util.LinkedList;
import java.util.List;
import java.util.UUID;
public class LoadSqlUtilV2 {
// show temporary partitions from ${db}.${table}
// alter table ${db}.${table} drop temporary partition $tp_name
public static String formatPath(String path, String nameFs) {
if (!path.startsWith("hdfs://")) {
if (path.startsWith("/")) {
path = "hdfs://" + nameFs + path;
} else {
path = "hdfs://" + nameFs + "/" + path;
}
}
if (!path.endsWith("*")) {
if (path.endsWith("/")) {
path = path + "*";
} else {
path = path + "/*";
}
}
return path;
}
public static String formatString(String sql) {
return " '" + sql + "' ";
}
public static String getLoadSql(String db, String table, String label, String path
, String format, String delimiter, String columns
, String where, String partition, String sets
, String brokerName, String hdfsUser, String nameFs, String namenode01, String namenode02
, String pro) {
StringBuilder sql = new StringBuilder();
sql.append("LOAD LABEL ").append(formatTable(db)).append(".").append(db).append("_").append(table).append("_").append(label).append("\n")
.append("(").append("\n")
.append("DATA INFILE('").append(formatPath(path, nameFs)).append("')").append("\n")
.append("INTO TABLE ").append(formatTable(table)).append("\n");
if (EmptyUtils.isNotEmpty(format)) {
sql.append("FORMAT AS ").append(format).append("\n");
} else if (EmptyUtils.isNotEmpty(delimiter)) {
sql.append("\nCOLUMNS TERMINATED BY ").append(LoadSqlUtilV2.formatString(delimiter)).append("\n");
} else {
System.out.println("导入失败,delimiter 和 format 不能同时为空,都不为空优先取 format ");
System.exit(-1);
}
sql.append("(").append(formatColumns(columns)).append(")");
if (EmptyUtils.isNotEmpty(where)) {
sql.append("\nwhere ").append(where);
}
if (EmptyUtils.isNotEmpty(partition)) {
sql.append("\nCOLUMNS FROM PATH AS (").append(partition).append(")");
}
if (EmptyUtils.isNotEmpty(sets)) {
sql.append("\nSET").append("(")
.append(sets)
.append(")");
}
sql.append("\n)").append("WITH BROKER ").append(LoadSqlUtilV2.formatString(brokerName)).append("\n");
if (EmptyUtils.isNotEmpty(nameFs)) {
sql.append("\n(").append("'username'=").append(LoadSqlUtilV2.formatString(hdfsUser)).append(",").append("\n")
.append("'dfs.nameservices'=").append(LoadSqlUtilV2.formatString(nameFs)).append(",").append("\n")
.append("'dfs.ha.namenodes.").append(nameFs).append("'='namenode01,namenode02',").append("\n")
.append("'dfs.namenode.rpc-address.qmai-cluster.namenode01'=").append(LoadSqlUtilV2.formatString(namenode01)).append(",").append("\n")
.append("'dfs.namenode.rpc-address.qmai-cluster.namenode02'=").append(LoadSqlUtilV2.formatString(namenode02)).append(",").append("\n")
.append("'dfs.client.failover.proxy.provider' = 'org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider'").append("\n")
.append(")");
}
// pro: timeout=3600
if (EmptyUtils.isNotEmpty(pro)) {
StringBuilder sbPro = new StringBuilder();
String[] split = pro.split(",");
for (String s : split) {
String[] split1 = s.split("=");
String key = split1[0].replaceAll("\"", "").replaceAll("'", "");
String value = split1[1].replaceAll("\"", "").replaceAll("'", "");
sbPro.append("'").append(key).append("'=").append("'").append(value).append("',");
}
sql.append("PROPERTIES").append("(").append(sbPro.substring(0, sbPro.length() - 1)).append(")");
}
return sql.toString();
}
public static String formatColumns(String columns) {
String[] split = columns
.replaceAll(" ", "")
.replaceAll("\n", "")
.replaceAll("'", "")
.replaceAll("`", "")
.split(",");
return "`" + StringUtils.join(Arrays.asList(split), "`,`") + "`";
}
public static String getUUID() {
return UUID.randomUUID().toString().replaceAll("-", "");
}
public static String showCreateTable(String db, String table) {
StringBuilder sql = new StringBuilder();
sql.append("show create table ").append(formatTable(db)).append(".").append(formatTable(table));
return sql.toString();
}
public static String deleteData(String db, String table, String dorisPartition, List<Partition> partitionValue) {
StringBuilder sql = new StringBuilder();
if (EmptyUtils.isNotEmpty(partitionValue) && partitionValue.size() > 0) {
sql.append("delete from ").append(formatTable(db)).append(".").append(formatTable(table)).append(" ")
.append("where ").append(dorisPartition).append("='").append(partitionValue.get(0).getValue()).append("'");
} else {
sql.append("TRUNCATE table ").append(formatTable(db)).append(".").append(formatTable(table));
}
return sql.toString();
}
public static List<Partition> getPartitionValue(String path) {
List<Partition> partitions = new LinkedList<>();
String[] split = path.split("/");
for (String s : split) {
if (s.contains("=")) {
String[] split1 = s.split("=");
if (split1.length == 2) {
partitions.add(new Partition(split1[0], split1[1]));
} else {
System.out.println("导入失败,hdfs 路径不正确:" + path);
System.exit(-1);
}
}
}
return partitions;
}
public static String formatTable(String table) {
return "`" + table.replaceAll(" ", "").replaceAll("`", "") + "`";
}
public static void main(String[] args) {
// String uuid = getUUID();
// String db = "dataviewdopro";
// String table = "dws_fina_order_items_di";
String path = "hdfs://qmai-cluster/user/hive/warehouse/dw_db.db/dws_fina_order_items_di/ds=2022-09-14";
// List<Partition> partitionValue = LoadSqlUtil.getPartitionValue(path);
// String co = "change_type,change_type_name,settle_type,settle_type_name,seller_id,store_id,statistics_date,business_date,date_hour,date_half_hour,reservation_order,order_source,order_type,pay_type,biz_type,table_area_id,table_area_name,table_id,table_name,meal_code,meal_name,event_type,item_level,item_sign,item_spec,is_gift,item_practice,item_type,sku_id,goods_id,name,backend_cate_id,backend_cate_name,num,goods_amount,receivable_amount,combined_spillover_amount,cost_amount,received_amount,ord_cnt,delivery_amt,pack_amt,tableware_amt,is_deleted,process_date,item_spec_id,item_practice_id";
// String loadSql = LoadSqlUtil.getLoadSql(db, table, uuid,
// path, "parquet", ""
// , co
// , "", "", partitionValue, "", "broker1", "hive", "qmai-cluster", "10.10.3.103:8020", "10.10.3.100:8020");
//
// System.out.println(loadSql);
//
// System.out.println(LoadSqlUtil.getLoadStatusSql(db, table, uuid));
//
// System.out.println(addPartition(db, table, "p", "", ""));
List<Partition> partitionValue = getPartitionValue(path);
System.out.println(partitionValue);
}
}
package com.qm.history;
import com.qm.doris.Load2DorisJdbc;
import com.qm.util.Constants;
import com.qm.util.DateUtils;
import com.qm.util.MysqlJdbcUtil;
import org.apache.commons.lang3.StringUtils;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.Date;
import java.util.LinkedList;
import java.util.List;
public class CopyTable {
public static void main(String[] args) {
history("2022-07-30", "2022-09-16", "canypro", "ads_fina_order_items_cate_di", 200000);
}
public static void history(String start1, String end1, String db, String table, Integer bath) {
Date startDate = new Date();
Date start = DateUtils.toDate(start1, "yyyy-MM-dd");
Date end = DateUtils.toDate(end1, "yyyy-MM-dd");
List<String> fields = descTable(db, table);
while (start.getTime() <= end.getTime()) {
String ds = DateUtils.toDateString(start, "yyyy-MM-dd");
System.out.println(ds + " 开始导入:" + DateUtils.toDateString(startDate, "yyyy-MM-dd HH:mm:ss"));
Long needDelete = getTotal("select count(*) from " + db + "." + table + "_tmp where process_date='" + ds + "' limit 1");
System.out.println(ds + " " + db + "." + table + "_tmp:需要删除 " + needDelete);
System.out.print(ds + " 正在删除... ");
for (int i = 0; i <= needDelete / bath; i++) {
System.out.print((i + 1) + ">> ");
Load2DorisJdbc.executeVoidSql("DELETE from " + db + "." + table + "_tmp where process_date='" + ds + "' limit " + bath, Constants.Tidb_url, Constants.Tidb_user, Constants.Tidb_password);
}
Long needLoad = getTotal("select count(*) from " + db + "." + table + " where process_date='" + ds + "' limit 1");
long m = needLoad / bath;
long n = needLoad % bath;
if (n > 0) {
m = m + 1;
}
System.out.println("\n" + ds + " " + db + "." + table + "_tmp:需要导入 " + needLoad + " 共:" + m + "批次");
System.out.print(ds + " 正在导入... ");
for (int i = 0; i < m; i++) {
System.out.print((i + 1) + ">> ");
Load2DorisJdbc.executeVoidSql("insert into " + db + "." + table + "_tmp (" + StringUtils.join(fields, ",") + ") "
+ "\nselect " + StringUtils.join(fields, ",") + " from " + db + "." + table
+ "\nwhere process_date='" + ds + "'"
+ "\norder by id"
+ "\nlimit " + (i * bath) + "," + bath, Constants.Tidb_url, Constants.Tidb_user, Constants.Tidb_password);
}
Date endDate = new Date();
System.out.println("\n" + ds + " 导入完成:" + DateUtils.toDateString(endDate, "yyyy-MM-dd HH:mm:ss") + " 共耗时:" + (endDate.getTime() - startDate.getTime()) + "");
start = DateUtils.addDate(start, 1);
}
}
public static Long getTotal(String sql) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(Constants.Tidb_url, Constants.Tidb_user, Constants.Tidb_password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
return res.getLong(1);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return 0L;
}
public static List<String> descTable(String db, String tableName) {
List<String> strings = new LinkedList<>();
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(Constants.Tidb_url, Constants.Tidb_user, Constants.Tidb_password);
stmt = conn.createStatement();
res = stmt.executeQuery("desc " + db + "." + tableName);
while (res.next()) {
strings.add(res.getString(1));
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return strings;
}
}
//package com.qm.history;
//
//import com.qm.doris.Load2Doris;
//import com.qm.doris.LoadSqlUtil;
//import com.qm.util.DateUtils;
//
//import java.util.Date;
//
//public class Insert2Doris {
// public static void main(String[] args) {
//
// String db = "dataviewdopro";
// String table = "dws_fina_order_items_di__tmp";
// String path = "/user/hive/warehouse/dw_db.db/dws_fina_order_items_di/ds";
// String format = "parquet";
// String delimiter = "";
// String columns = "change_type,change_type_name,settle_type,settle_type_name,seller_id,store_id,statistics_date,business_date,date_hour,date_half_hour,reservation_order,order_source,order_type,pay_type,biz_type,table_area_id,table_area_name,table_id,table_name,meal_code,meal_name,event_type,item_level,item_sign,item_spec,is_gift,item_practice,item_type,sku_id,goods_id,name,backend_cate_id,backend_cate_name,num,goods_amount,receivable_amount,combined_spillover_amount,cost_amount,received_amount,ord_cnt,delivery_amt,pack_amt,tableware_amt,is_deleted,process_date,item_spec_id,item_practice_id";
// String sets = "";
// String brokerName = "broker1";
// String hdfsUser = "hive";
// String nameFs = "qmai-cluster";
// String namenode01 = "10.10.3.103:8020";
// String namenode02 = "10.10.3.100:8020";
// String where = "";
// String url = "jdbc:mysql://47.100.132.12:9030/dataviewdopro?useSSL=false";
// String user = "bigdataviewdopro";
// String password = "bigdataviewdopro@123";
// String overwrite = "true";
// String partitionType = "day";
//
// Insert2Doris.history("2022-07-28", "2022-09-25", db, table, path, format, delimiter, columns, sets, overwrite, brokerName
// , hdfsUser, nameFs, namenode01, namenode02, where, url, user, password
// , partitionType);
//
// }
//
// public static void history(String start1, String end1, String db, String table, String path, String format, String delimiter
// , String columns, String sets, String overwrite
// , String brokerName, String hdfsUser, String nameFs, String namenode01, String namenode02
// , String where, String url, String user, String password
// , String partitionType) {
// Date startDate = new Date();
// Date start = DateUtils.toDate(start1, "yyyy-MM-dd");
// Date end = DateUtils.toDate(end1, "yyyy-MM-dd");
//
//// while (start.getTime() <= end.getTime()) {
//// String label = LoadSqlUtil.getUUID();
//// String ds = DateUtils.toDateString(start, "yyyy-MM-dd");
//// System.out.println(ds + " 开始导入:" + DateUtils.toDateString(startDate, "yyyy-MM-dd HH:mm:ss"));
//// Load2Doris.load2Doris(db, table, path + "=" + ds, format, delimiter, columns, sets, overwrite, brokerName
//// , hdfsUser, nameFs, namenode01, namenode02, where, url, user, password
//// , partitionType, label);
//// System.out.println("\n" + ds + " " + db + "." + table);
////
////
//// System.out.print(ds + " 正在导入... ");
////
//// Date endDate = new Date();
//// System.out.println("\n" + ds + " 导入完成:" + DateUtils.toDateString(endDate, "yyyy-MM-dd HH:mm:ss") + " 共耗时:" + (endDate.getTime() - startDate.getTime()) + "");
////
//// start = DateUtils.addDate(start, 1);
// }
// }
//
//}
package com.qm.history;
public class ItemCateDi {
public static void main(String[] args) {
CopyTable.history("2022-08-19", "2022-08-19", "canypro", "ads_fina_order_items_cate_di", 200000);
}
}
package com.qm.history;
public class ItemDi {
public static void main(String[] args) {
CopyTable.history("2022-08-28", "2022-09-15", "canypro", "ads_fina_order_items_di", 200000);
}
}
package com.qm.history;
import com.qm.doris.Load2DorisJdbc;
import com.qm.util.DateUtils;
import com.qm.util.MysqlJdbcUtil;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.Date;
public class XyTest {
public static void main(String[] args) {
String url = "jdbc:mysql://47.100.132.12:9030/dataviewdopro?useSSL=false";
String user = "dataviewdopro";
String password = "dataviewdopro@123";
XyTest.history("2022-07-28", "2022-09-26", url, user, password);
}
public static void history(String start1, String end1, String url, String user, String password) {
Date startDate = new Date();
Date start = DateUtils.toDate(start1, "yyyy-MM-dd");
Date end = DateUtils.toDate(end1, "yyyy-MM-dd");
while (start.getTime() <= end.getTime()) {
String ds = DateUtils.toDateString(start, "yyyy-MM-dd");
System.out.println(ds + " 开始导入:" + DateUtils.toDateString(startDate, "yyyy-MM-dd HH:mm:ss"));
Long needDelete = getTotal("select count(*) from dataviewdopro.dws_fina_order_items_di where process_date='" + ds + "' limit 1", url, user, password);
System.out.println(ds + " dataviewdopro.dws_fina_order_items_di :需要删除 " + needDelete);
System.out.print(ds + " 正在删除... ");
// for (int i = 0; i <= needDelete / 200000; i++) {
// System.out.print((i + 1) + ">> ");
Load2DorisJdbc.executeVoidSql("delete from dataviewdopro.dws_fina_order_items_di where process_date='" + ds + "' ", url, user, password);
// }
String sql = " insert into dataviewdopro.dws_fina_order_items_di\n" +
" (\n" +
" statistics_date\n" +
",date_hour\n" +
",date_half_hour\n" +
",seller_id\n" +
",store_id\n" +
",goods_id\n" +
",change_type_name\n" +
",settle_type_name\n" +
",process_date\n" +
",business_date\n" +
",change_type\n" +
",settle_type\n" +
",reservation_order\n" +
",order_source\n" +
",order_type\n" +
",biz_type\n" +
",table_area_id\n" +
",table_area_name\n" +
",table_id\n" +
",table_name\n" +
",meal_code\n" +
",meal_name\n" +
",event_type\n" +
",item_level\n" +
",item_sign\n" +
",item_spec\n" +
",is_gift\n" +
",item_practice\n" +
",item_type\n" +
",sku_id\n" +
",name\n" +
",backend_cate_id\n" +
",backend_cate_name\n" +
",is_deleted\n" +
",item_spec_id\n" +
",item_practice_id\n" +
",num\n" +
",goods_amount\n" +
",receivable_amount\n" +
",combined_spillover_amount\n" +
",cost_amount\n" +
",received_amount\n" +
",ord_cnt)\n" +
" select statistics_date\n" +
",date_hour\n" +
",date_half_hour\n" +
",seller_id\n" +
",store_id\n" +
",goods_id\n" +
",change_type_name\n" +
",settle_type_name\n" +
",process_date\n" +
",business_date\n" +
",change_type\n" +
",settle_type\n" +
",reservation_order\n" +
",order_source\n" +
",order_type\n" +
",biz_type\n" +
",table_area_id\n" +
",table_area_name\n" +
",table_id\n" +
",table_name\n" +
",meal_code\n" +
",meal_name\n" +
",event_type\n" +
",item_level\n" +
",item_sign\n" +
",item_spec\n" +
",is_gift\n" +
",item_practice\n" +
",item_type\n" +
",sku_id\n" +
",name\n" +
",backend_cate_id\n" +
",backend_cate_name\n" +
",is_deleted\n" +
",item_spec_id\n" +
",item_practice_id\n" +
",num\n" +
",goods_amount\n" +
",receivable_amount\n" +
",combined_spillover_amount\n" +
",cost_amount\n" +
",received_amount\n" +
",ord_cnt\n" +
"from dataviewdopro.dws_fina_order_items_di__tmp\n" +
"where process_date='" + ds + "'";
Load2DorisJdbc.executeVoidSql(sql, url, user, password);
System.out.print(ds + " 正在导入... ");
Date endDate = new Date();
System.out.println("\n" + ds + " 导入完成:" + DateUtils.toDateString(endDate, "yyyy-MM-dd HH:mm:ss") + " 共耗时:" + (endDate.getTime() - startDate.getTime()) + "");
start = DateUtils.addDate(start, 1);
}
}
public static Long getTotal(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
res = stmt.executeQuery(sql);
while (res.next()) {
return res.getLong(1);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return 0L;
}
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class BrokerLoadResult implements Serializable {
private String JobId;
private String Label;
private String State;//FINISHED
private String Progress;
private String Type;
private String EtlInfo;
private String TaskInfo;
private String ErrorMsg;
private String CreateTime;
private String EtlStartTime;
private String EtlFinishTime;
private String LoadStartTime;
private String LoadFinishTime;
private String URL;
private String JobDetails;
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class DescTable implements Serializable {
private String Field;
private String Type;
private String Null;
private String Key;
private String Default;
private String Extra;
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class Partition {
private String name;
private String value;
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
@AllArgsConstructor
@NoArgsConstructor
@Data
public class PartitionRange implements Serializable {
private String start;
private String end;
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class ShowCreateTable implements Serializable {
private String partition;
private Boolean dynamicPartition = false;
}
package com.qm.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class ShowPartitions implements Serializable {
private String PartitionId;
private String PartitionName;
private String VisibleVersion;
private String VisibleVersionTime;
private String VisibleVersionHash;
private String State;
private String PartitionKey;
private String Range; // 无分区或者非分区表的时候为空
private String DistributionKey;
private String Buckets;
private String ReplicationNum;
private String StorageMedium;
private String CooldownTime;
private String LastConsistencyCheckTime;
private String DataSize;
private String IsInMemory;
}
package com.qm.mysql;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
public class MysqlClient {
public static void main(String[] args) {
String url = args[0];
String user = args[1];
String password = args[2];
String sql = args[3];
MysqlClientUtil.executeSql(sql, url, user, password);
}
}
package com.qm.mysql;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONException;
import com.alibaba.fastjson.JSONObject;
import com.qm.util.MysqlJdbcUtil;
import java.sql.*;
import java.util.ArrayList;
import java.util.List;
public class MysqlClientUtil {
public static JSONArray executeSql(String sql, String url, String user, String password) {
Connection conn = null;
Statement stmt = null;
ResultSet res = null;
try {
conn = MysqlJdbcUtil.getConn(url, user, password);
stmt = conn.createStatement();
List<String> sqls = splitQueries(sql);
if (sqls.size() == 1) {
if (sql.trim().toLowerCase().startsWith("select")
|| sql.trim().toLowerCase().startsWith("show")
|| sql.trim().toLowerCase().startsWith("desc")
) {
res = stmt.executeQuery(sql);
resultSet2Array2(res);
} else {
stmt.executeUpdate(sql);
res = stmt.executeQuery("select 1");
}
} else {
for (int i = 0; i <= sqls.size() - 2; i++) {
stmt.executeUpdate(sqls.get(i));
}
String lastSql = sqls.get(sqls.size());
if (lastSql.trim().toLowerCase().startsWith("select")
|| lastSql.trim().toLowerCase().startsWith("show")
|| lastSql.trim().toLowerCase().startsWith("desc")
) {
res = stmt.executeQuery(lastSql);
resultSet2Array2(res);
} else {
stmt.executeUpdate(lastSql);
res = stmt.executeQuery("select 1");
}
}
} catch (Exception e) {
e.printStackTrace();
} finally {
MysqlJdbcUtil.closeAll(conn, stmt);
}
return null;
}
/**
* 按分号来切割SQL
*/
public static List<String> splitQueries(String queries) {
List<String> result = new ArrayList<>();
boolean inQuotes = false;
int start = 0;
for (int i = 0; i < queries.length(); i++) {
char c = queries.charAt(i);
if (c == '\'' || c == '"') {
inQuotes = !inQuotes;
} else if (!inQuotes && c == ';') {
result.add(queries.substring(start, i));
start = i + 1;
}
}
if (start < queries.length()) {
result.add(queries.substring(start));
}
return result;
}
public static JSONArray resultSet2Array(ResultSet res) throws SQLException, JSONException {
JSONArray array = new JSONArray();
ResultSetMetaData metaData = res.getMetaData();
int columnCount = metaData.getColumnCount();
while (res.next()) {
JSONObject json = new JSONObject();
for (int i = 1; i <= columnCount; i++) {
String columnName = metaData.getColumnLabel(i);
Object value = res.getObject(columnName);
json.put(columnName, value);
}
array.add(json);
}
return array;
}
public static void resultSet2Array2(ResultSet res) throws SQLException, JSONException {
ResultSetMetaData metaData = res.getMetaData();
int columnCount = metaData.getColumnCount();
while (res.next()) {
for (int i = 1; i <= columnCount; i++) {
String columnName = metaData.getColumnLabel(i);
Object value = res.getObject(columnName);
if (i == columnCount) {
System.out.print(value + " \n");
} else {
System.out.print(value + " ");
}
}
}
}
}
package com.qm.util;
import com.alibaba.fastjson.JSON;
import net.sf.cglib.beans.BeanCopier;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
/**
* 基于cglib进行Bean Copy
* @author 许育
*/
public class BeanUtil {
private BeanUtil() {}
/**
* 基于cglib进行对象复制
*
* @param source 被复制的对象
* @param clazz 复制对象类型
* @return
*/
public static <T> T copy(Object source, Class<T> clazz) {
if (EmptyUtils.isEmpty(source)) {
return null;
}
T target = instantiate(clazz);
BeanCopier copier = BeanCopier.create(source.getClass(), clazz, false);
copier.copy(source, target, null);
return target;
}
/**
* 基于cglib进行对象复制
*
* @param source 被复制的对象
* @param target 复制对象
* @return
*/
public static void copy(Object source, Object target) {
Objects.requireNonNull(source, "The source must not be null");
Objects.requireNonNull(target, "The target must not be null");
BeanCopier copier = BeanCopier.create(source.getClass(), target.getClass(), false);
copier.copy(source, target, null);
}
/**
* 基于cglib进行对象组复制
*
* @param datas 被复制的对象数组
* @param clazz 复制对象
* @return
*/
public static <T> List<T> copyByList(List<?> datas, Class<T> clazz) {
if (EmptyUtils.isEmpty(datas)) {
return Collections.emptyList();
}
List<T> result = new ArrayList<>(datas.size());
for (Object data : datas) {
result.add(copy(data, clazz));
}
return result;
}
/**
* 利用fastjson进行深拷贝
*
* @author delu.lv
* @param datas
* @param clazz
* @return
*/
public static <T> List<T> deepCopyByList(List<?> datas, Class<T> clazz) {
if (EmptyUtils.isEmpty(datas)) {
return Collections.emptyList();
}
return JSON.parseArray(JSON.toJSONString(datas), clazz);
}
/**
* 通过class实例化对象
*
* @param clazz
* @return
* @throws RuntimeException
*/
public static <T> T instantiate(Class<T> clazz) {
Objects.requireNonNull(clazz, "The class must not be null");
try {
return clazz.newInstance();
} catch (InstantiationException ex) {
throw new RuntimeException(clazz + ":Is it an abstract class?", ex);
} catch (IllegalAccessException ex) {
throw new RuntimeException(clazz + ":Is the constructor accessible?", ex);
}
}
}
package com.qm.util;
public class Constants {
public final static String Driver = "com.mysql.jdbc.Driver";
public final static String Tidb_url = "jdbc:mysql://47.100.132.12:4042/sqldataviewpro?characterEncoding=utf8&autoReconnect=true&failOverReadOnly=false&useSSL=false";
public final static String Tidb_user = "dwloadpro";
public final static String Tidb_password = "g$@2h0jk";
}
\ No newline at end of file
This diff is collapsed.
package com.qm.util;
import java.lang.reflect.Array;
import java.util.Collection;
import java.util.Map;
import java.util.Optional;
/**
* @author 许育
*/
public final class EmptyUtils {
private EmptyUtils() {
}
/**
* 判断一个对象是否为空
*
* @param obj
* @return
*/
public static boolean isEmpty(Object obj) {
if (obj == null) {
return true;
}
if (obj instanceof Optional) {
return !((Optional) obj).isPresent();
}
if (obj instanceof CharSequence) {
return ((CharSequence) obj).length() == 0;
}
if (obj.getClass().isArray()) {
return Array.getLength(obj) == 0;
}
if (obj instanceof Collection) {
return ((Collection) obj).isEmpty();
}
if (obj instanceof Map) {
return ((Map) obj).isEmpty();
}
// else
return false;
}
/**
* 判断一个对象是否不为空
*
* @param t
* @return
*/
public static boolean isNotEmpty(Object t) {
return !isEmpty(t);
}
/**
* 判断多个T是否存在空对象,只判断null不判断空 可用于多参数简化代码: 如: if (parameter1==null || parameter2==null || parameter3==null) 可以简化为: if
* (EmptyUtils.hasNull(parameter1, parameter2,parameter3))
*
* @param datas
* @return
* @author wangweizhen
*/
public static <T> boolean hasNull(T... datas) {
for (T t : datas) {
if (t == null) {
return true;
}
}
return false;
}
/**
* 判断多个Map是否存在空对象
*
* @param datas
* @param <K>
* @param <V>
* @return
*/
public static <K, V> boolean hasEmpty(Map<K, V>... datas) {
for (Map<K, V> data : datas) {
if (isEmpty(data)) {
return true;
}
}
return false;
}
/**
* 判断多个Collection是否存在空对象
*
* @param datas
* @param <T>
* @return
*/
public static <T> boolean hasEmpty(Collection<T>... datas) {
for (Collection<T> data : datas) {
if (isEmpty(data)) {
return true;
}
}
return false;
}
/**
* 找到第一个不为null的值
*
* @param data
* @param <T>
* @return
*/
public static <T> T ifNull(T... data) {
for (T d : data) {
if (d != null) {
return d;
}
}
return null;
}
public static void main(String[] args) {
System.out.println(EmptyUtils.isEmpty(""));
}
}
package com.qm.util;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import org.apache.http.HttpResponse;
import org.apache.http.client.HttpClient;
import org.apache.http.client.config.RequestConfig;
import org.apache.http.client.methods.HttpGet;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.client.methods.HttpPut;
import org.apache.http.client.methods.HttpUriRequest;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClientBuilder;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.message.BasicHeader;
import org.apache.http.util.EntityUtils;
import java.io.IOException;
import java.net.MalformedURLException;
import java.net.URI;
import java.net.URISyntaxException;
import java.net.URL;
import java.util.List;
import java.util.Map;
/**
* @author 许育
*/
public class HttpUtils {
public static String doJsonPost(String urlStr, Map<String, Object> data, Integer timeOut) {
HttpPost httpPost = new HttpPost(urlStr);
StringEntity stringEntity = new StringEntity(JSON.toJSONString(data), "UTF-8");
stringEntity.setContentType("application/json");
httpPost.setEntity(stringEntity);
return doExecute(urlStr, httpPost, timeOut);
}
public static String doJsonPostJson(String urlStr, JSON data, Integer timeOut) {
HttpPost httpPost = new HttpPost(urlStr);
StringEntity stringEntity = new StringEntity(data.toString(), "UTF-8");
stringEntity.setContentType("application/json");
httpPost.setEntity(stringEntity);
return doExecute(urlStr, httpPost, timeOut);
}
public static String doJsonPostJson(String urlStr, JSON data, List<BasicHeader> headers, Integer timeOut) {
HttpPost httpPost = new HttpPost(urlStr);
StringEntity stringEntity = new StringEntity(data.toString(), "UTF-8");
stringEntity.setContentType("application/json");
httpPost.setEntity(stringEntity);
if (headers != null && headers.size() > 0) {
for (BasicHeader header : headers) {
httpPost.setHeader(header);
}
}
return doExecute(urlStr, httpPost, timeOut);
}
public static String doJsonPut(String urlStr, Map<String, Object> data, Integer timeOut) {
HttpPut httpPut = new HttpPut(urlStr);
StringEntity stringEntity = new StringEntity(JSON.toJSONString(data), "UTF-8");
stringEntity.setContentType("application/json");
httpPut.setEntity(stringEntity);
return doExecute(urlStr, httpPut, timeOut);
}
public static String doGet(String urlStr, Integer timeOut) {
return doGet(urlStr, null, timeOut);
}
public static String doGet(String urlStr, List<BasicHeader> headers, Integer timeOut) {
URI uri = null;
try {
URL url = new URL(urlStr);
uri = new URI(url.getProtocol(), url.getHost() + ":" + url.getPort(), url.getPath(), url.getQuery(), null);
} catch (URISyntaxException | MalformedURLException e) {
System.out.println("url格式错误:" + e.getMessage());
}
HttpGet httpGet = new HttpGet(uri);
if (headers != null && headers.size() > 0) {
for (BasicHeader header : headers) {
httpGet.setHeader(header);
}
}
return doExecute(urlStr, httpGet, timeOut);
}
private static String doExecute(String url, HttpUriRequest request, Integer timeOut) {
try {
int timeout = EmptyUtils.isEmpty(timeOut) ? 30000 : timeOut;
RequestConfig config = RequestConfig.custom()
.setConnectTimeout(timeout)
.setConnectionRequestTimeout(timeout)
.setSocketTimeout(timeout).build();
CloseableHttpClient httpClient =
HttpClientBuilder.create().setDefaultRequestConfig(config).build();
HttpResponse response = httpClient.execute(request);
int code = response.getStatusLine().getStatusCode();
if (code == 200) {
return EntityUtils.toString(response.getEntity());
} else {
System.out.println(url + " http请求异常:" + response.getStatusLine().getStatusCode() + response.getEntity().toString());
}
} catch (IOException e) {
System.out.println("发送http请求失败:" + e.getMessage());
}
return null;
}
public static String getToken(String url) {
String res = "";
try {
HttpPost httpPost = new HttpPost();
HttpClient httpClient = HttpClients.createDefault();
httpPost.setHeader("content-type", "application/json");
httpPost.setURI(new URI(url));
HttpResponse response = httpClient.execute(httpPost);
String result = EntityUtils.toString(response.getEntity(), "utf-8");
JSONObject jsonObject = JSON.parseObject(result);
res = jsonObject.get("access_token").toString();
} catch (Exception e) {
return "query error:" + e.getMessage();
}
return res;
}
public static void main(String[] args) {
}
}
package com.qm.util;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.SQLException;
import java.sql.Statement;
public class MysqlJdbcUtil {
public static Connection getConn(String url, String user, String password) {
try {
Class.forName(Constants.Driver);
} catch (ClassNotFoundException e) {
e.printStackTrace();
}
Connection conn = null;
try {
conn = DriverManager.getConnection(url, user, password);
} catch (SQLException e) {
e.printStackTrace();
}
return conn;
}
public static void closeAll(Connection conn, Statement stmt) {
try {
if (stmt != null) stmt.close();
} catch (SQLException e) {
e.printStackTrace();
}
try {
if (conn != null) conn.close();
} catch (SQLException e) {
e.printStackTrace();
}
}
}
package com.qm.util;
public class QMUtil {
public static void sleep(Integer mins) {
try {
Thread.sleep(mins * 1000);
} catch (Exception e) {
}
}
}
package com.qm.util;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import java.sql.ResultSet;
import java.sql.ResultSetMetaData;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
/**
* @author 许育
*/
public class SqlUtil {
public static String sqlTrim(String sql) {
sql = sql.trim();
if (sql.endsWith(";")) {
sql = sql.substring(0, sql.length() - 1);
sql = sqlTrim(sql);
}
return sql;
}
/**
* 按分号来切割SQL
*/
public static List<String> splitQueries(String queries) {
List<String> result = new ArrayList<>();
boolean inQuotes = false;
int start = 0;
for (int i = 0; i < queries.length(); i++) {
char c = queries.charAt(i);
if (c == '\'' || c == '"') {
inQuotes = !inQuotes;
} else if (!inQuotes && c == ';') {
result.add(queries.substring(start, i));
start = i + 1;
}
}
if (start < queries.length()) {
result.add(queries.substring(start));
}
return result;
}
public static List<String> getSqlFormBrackets(String sql) {
List<String> result = new ArrayList<>();
int start = 1;
int end = 0;
for (int i = 0; i < sql.length(); i++) {
char c = sql.charAt(i);
if (c == '(') {
start = start + 1;
} else if (c == ')') {
end = end + 1;
}
if (start == end) {
result.add(sql.substring(0, i));
result.add(sql.substring(i + 1));
break;
}
}
return result;
}
public static JSONArray resultSet2Array(ResultSet res) throws SQLException {
JSONArray array = new JSONArray();
ResultSetMetaData metaData = res.getMetaData();
int columnCount = metaData.getColumnCount();
// 获取数据
while (res.next()) {
JSONObject json = new JSONObject();
for (int i = 1; i <= columnCount; i++) {
String columnName = metaData.getColumnLabel(i);
Object value = res.getObject(columnName);
json.put(columnName, value);
}
array.add(json);
}
return array;
}
}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment