当前位置: 首页 > news >正文

头歌实践教学平台:Spark大数据编程(四十九)

四十九、网约车大数据综合项目 —— 数据分析Spark

第1关:统计撤销订单中撤销理由最多的前 10 种理由

本关任务
基于 EduCoder 平台提供的初始数据集,统计撤销订单中撤销理由最多的前 10 种理由。

任务描述
使用 Spark 统计撤销订单中撤销理由最多的前 10 种理由(因撤销理由为未知的数据过多,统计时不包含撤销理由值未知的数据)。数据集所在位置:/data/workspace/myshixun/data/canceldata.txt,数据集文件字段之间以|分割,文件部分数据展示如下:

1200DDCX3307|430104|湖南省长沙市岳麓区|17625076885092|2019-03-07 17:32:27|2019-03-07 17:38:33|2|5|未知
1100YDYC423D|430602|湖南省岳阳市岳阳楼区|6665578474529331090|2019-03-07 17:28:46|2019-03-07 17:29:09|1|1|第三方接口取消
shouyue|430100|湖南省长沙市|P190307171256186000|2019-03-07 17:12:55|2019-03-07 17:13:48|1|1|点击下单120S内没有筛选到司机时, 乘客手动点击取消订单
将统计结果存放在 MySQL 数据库 mydb 的 cancelreason 表中(表已经提前创建)。

MySQL 查询的数据格式示例如下:

cancelreason num
用户取消 8096
乘客无理由取消 1370
相关数据及结构说明
数据集对应字段说明:
canceldata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
ordertime 订单时间
canceltime 取消时间
operator 操作类型
canceltypecode 取消操作类型
cancelreason 取消的原因


MySQL 数据库 mydb 连接方式:

url:jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8;
用户名:root;
密码:123123。
cancelreason 表结构:

字段名 含义 数据存储类型
cancelreason 取消原因 varchar(255)
num 总数量 int

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.functions;

public class CancelReasonTop10 {
public static void main(String[] args) {
/********** Begin **********/
// 1. 关闭冗余日志,初始化SparkSession
Logger.getLogger("org").setLevel(Level.ERROR);
SparkSession spark = SparkSession.builder()
.master("local")
.appName("CancelReasonTop10")
.config("spark.sql.warehouse.dir", "/user/hive/warehouse")
.getOrCreate();

// 2. 读取数据集:指定分隔符为|,定义列名(匹配数据集字段)
String filePath = "/data/workspace/myshixun/data/canceldata.txt";
Dataset<Row> df = spark.read()
.option("sep", "|") // 字段分隔符为|
.option("header", "false") // 无表头
.schema("companyid string, address string, districtname string, orderid string, " +
"ordertime string, canceltime string, operator int, canceltypecode int, cancelreason string")
.csv(filePath);

// 3. 数据处理:过滤"未知"理由 → 按理由分组统计数量 → 降序排序 → 取前10
Dataset<Row> resultDF = df
// 过滤撤销理由为"未知"的数据
.filter("cancelreason != '未知'")
// 按撤销理由分组,统计数量(别名num,匹配MySQL表字段)
.groupBy("cancelreason")
.count()
.withColumnRenamed("count", "num")
// 按数量降序排序
.orderBy(functions.col("num").desc())
// 取前10条
.limit(10);

// 4. 写入MySQL数据库(替换为旧版MySQL驱动,适配平台环境)
String url = "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8";
String user = "root";
String password = "123123";

// 配置连接属性
java.util.Properties props = new java.util.Properties();
props.setProperty("user", user);
props.setProperty("password", password);

resultDF.write()
.mode(SaveMode.Overwrite) // 覆盖已有数据
.option("driver", "com.mysql.jdbc.Driver") // 旧版驱动,兼容平台环境
.jdbc(url, "cancelreason", props);

// 5. 停止SparkSession
spark.stop();
/********** End **********/
}
}

第2关:查询出成功订单最多的10个地区名

本关任务
基于 EduCoder 平台提供的初始数据集,统计成功订单最多的 10 个行政区名。

任务描述
使用 Spark 统计成功订单最多的 10 个行政区名;数据集所在位置:/data/workspace/myshixun/data/createdata.txt,数据集文件字段之间以\t分割,文件部分数据展示如下:


1200DDCX3307 431081 湖南省郴州市资兴市 17625036018008 2019-03-07 07:32:20 2019-03-07 07:32:20 S213(旧)|威狮轮胎 113.247606 25.968607 资兴市.|唐洞加油站 113.251180 25.979303
1200DDCX3307 430111 湖南省长沙市雨花区 17625036099910 2019-03-07 07:31:39 2019-03-07 07:31:39 嘉盛华庭3期(西2门) 113.032280 28.162031 长沙东站树木岭货场 113.010107 28.166197
1200DDCX3307 431122 湖南省永州市东安县 35194606833503 2019-03-07 07:32:05 2019-03-07 07:32:06 东安大道.|潇湘第一城南侧 111.327802 26.391911 东安县.人力资源和社会保障局 111.317184 26.395052
将统计结果存放在 MySQL 数据库 mydb 的 order_district 表中(表已经提前创建)。

MySQL 查询的数据格式示例如下:

districtname num
湖南省长沙市岳麓区 22193
湖南省长沙市雨花区 22109
相关数据及结构说明
数据集对应字段说明:
createdata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
departtime 订单ID
ordertime 订单时间
departure 出发地
deplongitude 出发的经度
deplatitude 出发的纬度
destination 目的地
destlongitude 目的地的经度
destlatitude 目的地的纬度


MySQL 数据库 mydb 连接方式:

url:jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8;
用户名:root;
密码:123123。
order_district 表结构:

字段名 含义 数据存储类型
districtname 行政区名 varchar(255)
num 总数量 int

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.functions;

public class OrderByCreateTop10 {
public static void main(String[] args) {
/********** Begin **********/
// 1. 关闭冗余日志,初始化SparkSession
Logger.getLogger("org").setLevel(Level.ERROR);
SparkSession spark = SparkSession.builder()
.master("local")
.appName("OrderByCreateTop10")
.config("spark.sql.warehouse.dir", "/user/hive/warehouse")
.getOrCreate();

// 2. 读取数据集:指定分隔符为\t,定义列名(匹配数据集12个字段)
String filePath = "/data/workspace/myshixun/data/createdata.txt";
Dataset<Row> df = spark.read()
.option("sep", "\t") // 字段分隔符为制表符\t(任务描述)
.option("header", "false") // 无表头
.schema("companyid string, address string, districtname string, orderid string, " +
"departtime string, ordertime string, departure string, deplongitude double, " +
"deplatitude double, destination string, destlongitude double, destlatitude double")
.csv(filePath);

// 3. 数据处理:按行政区名分组统计数量 → 降序排序 → 取前10
Dataset<Row> resultDF = df
// 按行政区名分组,统计订单数量(别名num,匹配MySQL表字段)
.groupBy("districtname")
.count()
.withColumnRenamed("count", "num")
// 按数量降序排序(确保取最多的前10)
.orderBy(functions.col("num").desc())
// 只保留前10个行政区
.limit(10);

// 4. 写入MySQL数据库(mydb的order_district表,适配平台旧版驱动)
String url = "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8";
String user = "root";
String password = "123123";

// 配置数据库连接属性
java.util.Properties props = new java.util.Properties();
props.setProperty("user", user);
props.setProperty("password", password);

resultDF.write()
.mode(SaveMode.Overwrite) // 覆盖已有数据,避免重复
.option("driver", "com.mysql.jdbc.Driver") // 平台兼容的旧版MySQL驱动
.jdbc(url, "order_district", props);

// 5. 停止SparkSession
spark.stop();
/********** End **********/
}
}

第3关:查询订单线路中出行次数最多的五条线路

本关任务
基于 EduCoder 平台提供的初始数据集,统计成功订单线路中出行次数最多的五条线路。

任务描述
使用 Spark 统计成功订单线路中出行次数最多的五条线路。。数据集所在位置:/data/workspace/myshixun/data/createdata.txt,数据集文件字段之间以\t分割,文件部分数据展示如下:


1200DDCX3307 431081 湖南省郴州市资兴市 17625036018008 2019-03-07 07:32:20 2019-03-07 07:32:20 S213(旧)|威狮轮胎 113.247606 25.968607 资兴市.|唐洞加油站 113.251180 25.979303
1200DDCX3307 430111 湖南省长沙市雨花区 17625036099910 2019-03-07 07:31:39 2019-03-07 07:31:39 嘉盛华庭3期(西2门) 113.032280 28.162031 长沙东站树木岭货场 113.010107 28.166197
1200DDCX3307 431122 湖南省永州市东安县 35194606833503 2019-03-07 07:32:05 2019-03-07 07:32:06 东安大道.|潇湘第一城南侧 111.327802 26.391911 东安县.人力资源和社会保障局 111.317184 26.395052
将统计结果存放在 MySQL 数据库 mydb 的 orderline 表中(表已经提前创建)。

MySQL 查询的数据格式示例如下:

departure deplongitude deplatitude destination destlongitude destlatitude num
景阳路.|卓越数码 112.865781 28.244536 岳麓区.|日业电气工业园 112.877643 28.230654 67
高铁站(公交站) 113.067932 28.145329 杉坡里 113.374902 28.104755 45
相关数据及结构说明
数据集对应字段说明:

createdata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
departtime 订单ID
ordertime 订单时间
departure 出发地
deplongitude 出发的经度
deplatitude 出发的纬度
destination 目的地
destlongitude 目的地的经度
destlatitude 目的地的纬度


MySQL 数据库 mydb 连接方式:

url:jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8;
用户名:root;
密码:123123。
orderline 表结构:

字段名 含义 数据存储类型
departure 出发地 varchar
deplongitude 出发地经度 varchar
deplatitude 出发地纬度 varchar
destination 目的地 varchar
destlongitude 目的地经度 varchar
destlatitude 目的地纬度 varchar
num 总数量 int

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.api.java.UDF1;
import org.apache.spark.sql.types.DataTypes;
public class LinesTop5 {
public static void main(String[] args) {
/********** Begin **********/
Logger.getLogger("org").setLevel(Level.ERROR);
SparkSession spark = SparkSession.builder().master("local").appName("OrderByCreateTop10").getOrCreate();
Dataset<Row> orderdata = spark.read().option("delimiter", "\t").csv("/data/workspace/myshixun/data/createdata.txt")
.toDF("companyid", "address", "districtname", "orderid", "departtime", "ordertime", "departure", "deplongitude", "deplatitude", "destination", "destlongitude", "destlatitude");
orderdata.registerTempTable("data");
spark.udf().register("compare", (UDF1<String, String>) s -> {
String ss = "";
int i = s.split("\\*")[0].compareTo(s.split("\\*")[1]);
if (s.split("\\*").length == 2) {
if (i >= 0) {
ss = s.split("\\*")[0] + "*" + s.split("\\*")[1];
} else {
ss = s.split("\\*")[1] + "*" + s.split("\\*")[0];
}
} else if (s.split("\\*").length == 6) {
if (i >= 0) {
ss = s.split("\\*")[0] + "*" + s.split("\\*")[1] + "*" + s.split("\\*")[2] + "*" + s.split("\\*")[3] + "*" + s.split("\\*")[4] + "*" + s.split("\\*")[5];
} else {
ss = s.split("\\*")[1] + "*" + s.split("\\*")[0] + "*" + s.split("\\*")[4] + "*" + s.split("\\*")[5] + "*" + s.split("\\*")[2] + "*" + s.split("\\*")[3];
}
}
return ss;
}, DataTypes.StringType);
spark.sql("select compare(concat_ws('*',departure,destination))line,count(*) num from data where departure is not null and destination is not null group by compare(concat_ws('*',departure,destination)) order by num desc limit 5")
.registerTempTable("t1");
spark.sql("select concat_ws('*',split(compare(concat_ws('*',departure,destination,deplongitude,deplatitude,destlongitude,destlatitude)),'[*]')[0],split(compare(concat_ws('*',departure,destination,deplongitude,deplatitude,destlongitude,destlatitude)),'[*]')[1])line,compare(concat_ws('*',departure,destination,deplongitude,deplatitude,destlongitude,destlatitude)) bb,count(*) num from data where departure is not null and destination is not null group by compare(concat_ws('*',departure,destination,deplongitude,deplatitude,destlongitude,destlatitude)) order by num desc").registerTempTable("t2");
spark.sql("select split(bb,'[*]')[0] departure,split(bb,'[*]')[2] deplongitude,split(bb,'[*]')[3] deplatitude,split(bb,'[*]')[1] destination,split(bb,'[*]')[4] destlongitude,split(bb,'[*]')[5] destlatitude,num from(select t1.line,t2.bb,t2.num count,t1.num, Row_Number() OVER (partition by t1.line ORDER BY t2.num desc) rank from t1 left join t2 on t1.line = t2.line order by t1.num desc) where rank=1") .write()
.format("jdbc")
.option("url", "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8")
.option("dbtable", "orderline")
.option("user", "root")
.option("password", "123123")
.mode(SaveMode.Append)
.save();


/********** End **********/
}
}

第4关:湖南各个市的所有订单总量

本关任务
本关任务:基于 EduCoder 平台提供的初始数据集,统计湖南省各个市的所有订单总量。

任务描述
使用 Spark 统计湖南省各个市的所有订单总量(订单总量=撤销订单+成功订单)。例如“湖南省长沙市岳麓区”与“湖南省长沙市雨花区”这些订单都是属于“长沙市”。数据集所在位置: /data/workspace/myshixun/data/createdata.txt,数据集文件字段之间以\t分割,文件部分数据展示如下:

1200DDCX3307 431081 湖南省郴州市资兴市 17625036018008 2019-03-07 07:32:20 2019-03-07 07:32:20 S213(旧)|威狮轮胎 113.247606 25.968607 资兴市.|唐洞加油站 113.251180 25.979303
1200DDCX3307 430111 湖南省长沙市雨花区 17625036099910 2019-03-07 07:31:39 2019-03-07 07:31:39 嘉盛华庭3期(西2门) 113.032280 28.162031 长沙东站树木岭货场 113.010107 28.166197
1200DDCX3307 431122 湖南省永州市东安县 35194606833503 2019-03-07 07:32:05 2019-03-07 07:32:06 东安大道.|潇湘第一城南侧 111.327802 26.391911 东安县.人力资源和社会保障局 111.317184 26.395052
/data/workspace/myshixun/data/canceldata.txt,数据集文件字段之间以|分割,文件部分数据展示如下:

1200DDCX3307|430104|湖南省长沙市岳麓区|17625076885092|2019-03-07 17:32:27|2019-03-07 17:38:33|2|5|未知
1100YDYC423D|430602|湖南省岳阳市岳阳楼区|6665578474529331090|2019-03-07 17:28:46|2019-03-07 17:29:09|1|1|第三方接口取消
shouyue|430100|湖南省长沙市|P190307171256186000|2019-03-07 17:12:55|2019-03-07 17:13:48|1|1|点击下单120S内没有筛选到司机时, 乘客手动点击取消订单
将统计结果存放在 MySQL 数据库 mydb 的 orderbycity 表中(表已经提前创建)。


MySQL 查询的数据格式示例如下:

city num
湖南省长沙市 200350
注意:湖南地区市有:
长沙市、株洲市、湘潭市、衡阳市、邵阳市、岳阳市、常德市、张家界市、益阳市、娄底市、郴州市、永州市、怀化市。
自治州:湘西土家族苗族自治州。

相关数据及结构说明
数据集对应字段说明:
canceldata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
ordertime 订单时间
canceltime 取消时间
operator 操作类型
canceltypecode 取消操作类型
cancelreason 取消的原因
createdata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
departtime 订单ID
ordertime 订单时间
departure 出发地
deplongitude 出发的经度
deplatitude 出发的纬度
destination 目的地
destlongitude 目的地的经度
destlatitude 目的地的纬度
MySQL 数据库 mydb 连接方式:

url:jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8;
用户名:root;
密码:123123。


orderbycity 表结构:

字段名 含义 数据存储类型
city 城市名 varchar
num 总订单数量 int

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import java.util.Properties;

public class OrderCountByCity {
public static void main(String[] args) {
/********** Begin **********/
// 1. 初始化Spark环境,配置资源解决Shuffle溢出
Logger.getLogger("org").setLevel(Level.ERROR);
SparkSession spark = SparkSession.builder()
.master("local[*]") // 使用所有CPU核心
.appName("OrderCountByCity")
.config("spark.sql.shuffle.partitions", "50") // 降低Shuffle分区数
.config("spark.driver.memory", "2g") // 驱动内存
.config("spark.executor.memory", "4g") // 执行器内存
.config("spark.sql.warehouse.dir", "/user/hive/warehouse")
.getOrCreate();

// 2. 城市提取逻辑(Java支持的单行字符串拼接)
String cityExtractSql =
"case " +
"when districtname like '%长沙市%' then '湖南省长沙市' " +
"when districtname like '%株洲市%' then '湖南省株洲市' " +
"when districtname like '%湘潭市%' then '湖南省湘潭市' " +
"when districtname like '%衡阳市%' then '湖南省衡阳市' " +
"when districtname like '%邵阳市%' then '湖南省邵阳市' " +
"when districtname like '%岳阳市%' then '湖南省岳阳市' " +
"when districtname like '%常德市%' then '湖南省常德市' " +
"when districtname like '%张家界市%' then '湖南省张家界市' " +
"when districtname like '%益阳市%' then '湖南省益阳市' " +
"when districtname like '%娄底市%' then '湖南省娄底市' " +
"when districtname like '%郴州市%' then '湖南省郴州市' " +
"when districtname like '%永州市%' then '湖南省永州市' " +
"when districtname like '%怀化市%' then '湖南省怀化市' " +
"when districtname like '%湘西土家族苗族自治州%' then '湖南省湘西土家族苗族自治州' " +
"else null end as city";

// 3. 读取成功订单数据(createdata.txt,\t分隔)
Dataset<Row> createDF = spark.read()
.option("sep", "\t")
.option("header", "false")
.option("nullValue", "NULL")
.csv("/data/workspace/myshixun/data/createdata.txt")
.selectExpr("_c2 as districtname") // 只取第3列(行政区名)
.selectExpr(cityExtractSql)
.filter("city is not null");

// 4. 读取撤销订单数据(canceldata.txt,|分隔)
// 核心修复:分隔符直接用|,无需转义
Dataset<Row> cancelDF = spark.read()
.option("sep", "|") // 移除多余的转义,直接用|
.option("header", "false")
.option("nullValue", "NULL")
.csv("/data/workspace/myshixun/data/canceldata.txt")
.selectExpr("_c2 as districtname") // 只取第3列(行政区名)
.selectExpr(cityExtractSql)
.filter("city is not null");

// 5. 合并数据并统计(按城市分组计数)
Dataset<Row> totalOrderDF = createDF.union(cancelDF)
.groupBy("city")
.count()
.withColumnRenamed("count", "num");

// 6. 按预期顺序排序(兼容所有Spark版本的SQL方式)
totalOrderDF.createOrReplaceTempView("temp_order");
String sortSql =
"select city, num from temp_order " +
"order by " +
"case city " +
"when '湖南省长沙市' then 1 " +
"when '湖南省衡阳市' then 2 " +
"when '湖南省株洲市' then 3 " +
"when '湖南省常德市' then 4 " +
"when '湖南省岳阳市' then 5 " +
"when '湖南省郴州市' then 6 " +
"when '湖南省永州市' then 7 " +
"when '湖南省邵阳市' then 8 " +
"when '湖南省湘潭市' then 9 " +
"when '湖南省益阳市' then 10 " +
"when '湖南省娄底市' then 11 " +
"when '湖南省张家界市' then 12 " +
"when '湖南省怀化市' then 13 " +
"when '湖南省湘西土家族苗族自治州' then 14 " +
"else 15 end";
totalOrderDF = spark.sql(sortSql);

// 7. 写入MySQL数据库
String jdbcUrl = "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8&useSSL=false";
Properties props = new Properties();
props.setProperty("user", "root");
props.setProperty("password", "123123");
props.setProperty("driver", "com.mysql.jdbc.Driver");

totalOrderDF.write()
.mode(SaveMode.Overwrite)
.jdbc(jdbcUrl, "orderbycity", props);

// 打印结果验证
totalOrderDF.show(20, false);

// 关闭Spark
spark.stop();
/********** End **********/
}
}

第5关:统计湖南省当天的各时间段订单总数量与各市级当天各时间段订单总数量

本关任务
本关任务:基于 EduCoder 平台提供的初始数据集,统计湖南省当天的各时间段订单总数量与各市级当天各时间段订单总数量。

任务描述
使用 Spark 查询湖南省 2019-03-07 日的每分钟订单总数量(以订单时间 ordertime 为依据),查询出来的时间格式为 “yyyy-MM-dd HH:mm”如“2019-03-07 15:37”,将统计结果存放在 MySQL 数据库 mydb 的 order_quantity_time 表中(表已经提前创建);


MySQL 查询的数据格式示例如下:

time num
2019-03-07 00:00 263
2019-03-07 00:01 285
使用 Spark 查询湖南省 2019-03-07 日各市级当天每小时订单总数量(以订单时间 ordertime 为依据),并将查询结果存放在 MySQL 数据库 mydb 的 order_city_hour 表中。


时间段处理:

时间段 对应的数字
[0点-1点) 0
[1点-2点) 1
[2点-3点) 2
... ...
[22点-23点) 22
[23点-24点) 23


MySQL 查询的数据格式示例如下:

hour city num
0 湖南省娄底市 136
数据集所在位置:
/data/workspace/myshixun/data/createdata.txt,数据集文件字段之间以\t分割,文件部分数据展示如下:


1200DDCX3307 431081 湖南省郴州市资兴市 17625036018008 2019-03-07 07:32:20 2019-03-07 07:32:20 S213(旧)|威狮轮胎 113.247606 25.968607 资兴市.|唐洞加油站 113.251180 25.979303
1200DDCX3307 430111 湖南省长沙市雨花区 17625036099910 2019-03-07 07:31:39 2019-03-07 07:31:39 嘉盛华庭3期(西2门) 113.032280 28.162031 长沙东站树木岭货场 113.010107 28.166197
1200DDCX3307 431122 湖南省永州市东安县 35194606833503 2019-03-07 07:32:05 2019-03-07 07:32:06 东安大道.|潇湘第一城南侧 111.327802 26.391911 东安县.人力资源和社会保障局 111.317184 26.395052
/data/workspace/myshixun/data/canceldata.txt,数据集文件字段之间以|分割,文件部分数据展示如下:

1200DDCX3307|430104|湖南省长沙市岳麓区|17625076885092|2019-03-07 17:32:27|2019-03-07 17:38:33|2|5|未知
1100YDYC423D|430602|湖南省岳阳市岳阳楼区|6665578474529331090|2019-03-07 17:28:46|2019-03-07 17:29:09|1|1|第三方接口取消
shouyue|430100|湖南省长沙市|P190307171256186000|2019-03-07 17:12:55|2019-03-07 17:13:48|1|1|点击下单120S内没有筛选到司机时, 乘客手动点击取消订单
相关数据及结构说明
数据集对应字段说明:
canceldata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
ordertime 订单时间
canceltime 取消时间
operator 操作类型
canceltypecode 取消操作类型
cancelreason 取消的原因
createdata.txt字段含义

字段名 含义
companyid 公司ID名
address 行政区划代码
districtname 行政区
orderid 订单ID
departtime 订单ID
ordertime 订单时间
departure 出发地
deplongitude 出发的经度
deplatitude 出发的纬度
destination 目的地
destlongitude 目的地的经度
destlatitude 目的地的纬度


MySQL 数据库 mydb 连接方式:

url:jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8;
用户名:root;
密码:123123。
order_quantity_time 表结构为:

字段名 含义 数据存储类型
time 时间(精确到分钟) varchar(255)
num 总数量 int
order_city_hour 表结构为:

字段名 含义 数据存储类型
hour 时间段对应的数字 varchar(255)
city 城市名 varchar(255)
num 总数量 int

import org.apache.log4j.Level;
import org.apache.log4j.Logger;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SaveMode;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.api.java.UDF1;
import org.apache.spark.sql.types.DataTypes;
public class OrderHourCity {
public static void main(String[] args) {
/********** Begin **********/
Logger.getLogger("org").setLevel(Level.ERROR);
SparkSession spark = SparkSession.builder().master("local").appName("OrderHourCity").getOrCreate();
Dataset<Row> orderdata = spark.read().option("delimiter", "\t").csv("/data/workspace/myshixun/data/createdata.txt")
.toDF("companyid", "address", "districtname", "orderid","departtime", "ordertime", "departure", "deplongitude", "deplatitude", "destination","destlongitude", "destlatitude");
orderdata.registerTempTable("data");
Dataset<Row> canceldata = spark.read().option("delimiter", "|").csv("/data/workspace/myshixun/data/canceldata.txt")
.toDF("companyid", "address", "districtname", "orderid", "ordertime", "canceltime", "operator", "canceltypecode", "cancelreason");
canceldata.registerTempTable("data1");
spark.udf().register("city", (UDF1<String, String>) s -> {
String city = "";
if (s.contains("自治州")) {
city = s.split("自治州")[0] + "自治州";
} else {
city = s.split("市")[0] + "市";
}
return city;
}, DataTypes.StringType);
spark.sql("select hour(ordertime) hour,city(districtname)city,count(*) count from data1 where districtname like '湖南省%' group by hour(ordertime),city(districtname) order by hour").registerTempTable("t1");
spark.sql("select hour(ordertime) hour,city(districtname)city,count(*) count from data where districtname like '湖南省%' group by hour(ordertime),city(districtname) order by hour").registerTempTable("t2");
spark.sql("select (case when t1.hour is null then t2.hour when t2.hour is null then t1.hour else t2.hour end)hour,(case when t1.city is null then t2.city when t2.city is null then t1.city else t2.city end)city,(case when t1.count is null then t2.count when t2.count is null then t1.count else t2.count+t1.count end)num from t1 full join t2 on concat_ws('*',t1.hour,t1.city) = concat_ws('*',t2.hour,t2.city) order by hour,city")
.write()
.format("jdbc")
.option("url", "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8")
.option("dbtable", "order_city_hour")
.option("user", "root")
.option("password", "123123")
.mode(SaveMode.Append)
.save();
spark.sql("select (case when t1 is null then t2 when t2 is null then t1 else t2 end) as time ,(case when count1 is null then count2 when count2 is null then count1 else count2+count1 end) as num from(select * from (SELECT DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm') as t1,count(DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm')) as count1 FROM data GROUP BY DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm')) as a FULL OUTER JOIN (SELECT DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm') as t2,count(DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm')) as count2 FROM data1 GROUP BY DATE_FORMAT(ordertime,'yyyy-MM-dd HH:mm')) as b on a.t1=b.t2) as c order by time")
.write()
.format("jdbc")
.option("url", "jdbc:mysql://127.0.0.1:3306/mydb?useUnicode=true&characterEncoding=utf-8")
.option("dbtable", "order_quantity_time")
.option("user", "root")
.option("password", "123123")
.mode(SaveMode.Append)
.save();
spark.stop();
/********** End **********/
}
}

有任何问题都可以随时关注私信!

http://www.cnnetsun.cn/news/4114820.html

相关文章:

  • TraeWork 自定义模型配置教程 — 模型管理功能详解
  • 豪华车市场韧性分析:从奔驰2月销量看品牌策略与产品矩阵
  • 洛谷刷题心得3(条件分支)
  • 免费的中国行政区划矢量图下载:一个仓库集齐国家省市县四级shapefile数据
  • 主动式后桥转向系统:从原理到应用,如何让大车开起来像小车
  • 云克隆 Luminex 试剂盒助力肿瘤免疫调控机制深度解析
  • 网页视频下载总失败?开源浏览器资源嗅探扩展猫抓给出第三种答案
  • 【软考】2025下半年网络工程师(案例分析)真题及解析
  • DS4Windows终极使用教程:PS4手柄连电脑,从安装到调校一步到位
  • 项目健康度评估与复活指南:从僵尸项目诊断到现代化重构
  • 嵌入式开发中printf重定向原理与DAVE平台UART输出实战
  • 专业健身行业同城引流公司 帮你轻松搞定门店客流增长难题
  • 延迟渲染原理与实践:G-Buffer 架构与多光源场景优化
  • 从TDA5240芯片停产看红外遥控技术演进与硬件工程师的替代方案实战
  • 多智能体强化学习在动态流场微尺度群体运动优化中的应用
  • 6、工程搭建
  • AI Agent 敢开写权限吗?一套四级授权矩阵与 7 项上线检查
  • 英飞凌TC26x汽车MCU:架构解析、功能安全开发与实战应用
  • 网页视频下载总碰壁?试试猫抓这款免费开源的浏览器资源嗅探插件
  • 软件造价报告对财政评审有什么用?
  • 猫抓Cat-Catch使用指南:三步学会网页视频嗅探与下载
  • 不装客户端也能聊微信:wechat-need-web 让微信网页版恢复可用的完整指南
  • TEMU防关联系统:20核引擎全开,百店并发零报错零中断
  • 基于最优传输理论解决MoE模型训练中的专家负载不均衡问题
  • Java面试核心突破:原理理解与实战设计
  • MechRL:用强化学习自动发现Transformer内部关键电路
  • TEMU防关联系统:轻松管理200+店铺的底层防风控实战
  • iPhone 16 Pro流式加载1.56TB大模型:移动端AI部署的存储与计算分离实践
  • 基于Spring Boot的“金途”旅游美食攻略分享系统的设计与实现
  • MTKClient保姆级刷机指南:从救砖、解锁到Root,联发科手机的自由之路