赞
踩
今天的任务是完成 DWD 层剩余的事实表;今年的秋招开得比往年早,所以要抓紧时间了,据了解,今年的 hc 还是不多,要是晚点投铁定寄中寄了;
今天还是个周末,不过记忆里我好像整个大学都没有好好放松过过一个周末,毕竟只有周末的空教室才是最多的,可以抢到带插座的位置;emm... 等一切尘埃落定一定好好放松放松~
昨天我们创建了订单与处理表(把订单明细、订单、订单明细活动、订单明细优惠券等订单相关和表和 MySQL lookup 字典表关联到了一起),就是为了之后做关于订单的事实表的时候不用再去频繁关联造成重复计算;
这里只需要保留有用的数据,其它字段(比如 old 、order_status 这些字段这里完全没用,是给取消订单表用的)直接丢弃掉
- // TODO 2. 读取订单预处理表
- tableEnv.executeSql("create table dwd_trade_order_pre_process(" +
- "id ," +
- "order_id ," +
- "user_id ," +
- "sku_id ," +
- "sku_name ," +
- "province_id ," +
- "activity_id ," +
- "activity_rule_id ," +
- "coupon_id ," +
- "date_id ," +
- "create_time ," +
- "operate_date_id ," +
- "operate_time ," +
- "source_id ," +
- "source_type_id ," +
- "source_type_name ," +
- "sku_num ," +
- "split_original_amount ," +
- "split_activity_amount ," +
- "split_coupon_amount ," +
- "split_total_amount ," +
- "`type`" +
- ")" + MyKafkaUtil.getKafkaDDL("dwd_trade_order_pre_process", "trade_detail"));
订单预处理表的粒度是商品,所以一旦有订单被下单,这里新增记录数 = 订单内商品件数
- // TODO 3. 过滤出下单数据,即新增数据
- Table filterTable = tableEnv.sqlQuery("SELECT " +
- "id ," +
- "order_id ," +
- "user_id ," +
- "sku_id ," +
- "sku_name ," +
- "sku_num ," +
- "province_id ," +
- "activity_id ," +
- "activity_rule_id ," +
- "coupon_id ," +
- "create_time ," +
- "operate_date_id ," +
- "operate_time ," +
- "source_id ," +
- "source_type_id ," +
- "source_type_name ," +
- "split_activity_amount ," +
- "split_coupon_amount ," +
- "split_total_amount" +
- ") " +
- "FROM dwd_trade_order_pre_process " +
- "WHERE `type`='insert'"
- );
- tableEnv.createTemporaryView("filter_table",filterTable);
- // TODO 4. 创建 dwd 层下单数据表
- tableEnv.executeSql("CREATE TABLE dwd_trade_order_detail (" +
- " id STRING," +
- " order_id STRING," +
- " user_id STRING," +
- " sku_id STRING," +
- " sku_name STRING," +
- " sku_num STRING," +
- " province_id STRING," +
- " activity_id STRING," +
- " activity_rule_id STRING," +
- " coupon_id STRING," +
- " create_time STRING," +
- " operate_date_id STRING," +
- " operate_time STRING," +
- " source_id STRING," +
- " source_type_id STRING," +
- " source_type_name STRING," +
- " split_activity_amount STRING," +
- " split_coupon_amount STRING," +
- " split_total_amount STRING" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_trade_order_detail"));
- // TODO 5. 将数据写出到 kafka
- tableEnv.executeSql("INSERT INTO dwd_trade_order_detail SELECT * FROM filter_table");
这个需求也非常简单
思路:
- public class DwdTradeCancelDetail {
- public static void main(String[] args) throws Exception {
- // TODO 1. 获取执行环境
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(1); // 生产环境中设置为kafka主题的分区数
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // 1.1 开启checkpoint
- env.enableCheckpointing(5 * 60000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointStorage("hdfs://hadoop102:8020/s/ck");
- env.getCheckpointConfig().setCheckpointTimeout(10 * 60000L);
- env.getCheckpointConfig().setMaxConcurrentCheckpoints(2); // 设置最大共存的checkpoint数量
- env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3,5000L)); // 固定频率重启: 尝试3次重启,每5s重启一次
-
- // 1.2 设置状态后端
- env.setStateBackend(new HashMapStateBackend());
-
- // TODO 2. 读取订单预处理表
- tableEnv.executeSql("create table dwd_trade_order_pre_process(" +
- "id ," +
- "order_id ," +
- "user_id ," +
- "order_status ," +
- "sku_id ," +
- "sku_name ," +
- "province_id ," +
- "activity_id ," +
- "activity_rule_id ," +
- "coupon_id ," +
- "date_id ," +
- "create_time ," +
- "operate_date_id ," +
- "operate_time ," +
- "source_id ," +
- "source_type_id ," +
- "source_type_name ," +
- "sku_num ," +
- "split_original_amount ," +
- "split_activity_amount ," +
- "split_coupon_amount ," +
- "split_total_amount ," +
- "`type` ," +
- "`old` " +
- ")" + MyKafkaUtil.getKafkaDDL("dwd_trade_order_pre_process", "trade_detail"));
-
- // TODO 3. 过滤出取消订单的数据
- Table filterTable = tableEnv.sqlQuery("" +
- "SELECT" +
- " id ," +
- " order_id ," +
- " user_id ," +
- " order_status ," +
- " sku_id ," +
- " sku_name ," +
- " province_id ," +
- " activity_id ," +
- " activity_rule_id ," +
- " coupon_id ," +
- " date_id ," +
- " create_time ," +
- " operate_date_id ," +
- " operate_time ," +
- " source_id ," +
- " source_type_id ," +
- " source_type_name ," +
- " sku_num ," +
- " split_original_amount ," +
- " split_activity_amount ," +
- " split_coupon_amount ," +
- " split_total_amount ," +
- " `type` ," +
- " `old` " +
- "FROM dwd_trade_order_pre_process " +
- "WHERE `type`='update' " +
- "old['order_status'] is not null" +
- "AND order_status='1003'"
- );
- tableEnv.createTemporaryView("filter_table",filterTable);
-
- // TODO 4. 创建 dwd 层取消订单表
- tableEnv.executeSql("CREATE TABLE dwd_trade_cancel_detail" +
- "id STRING," +
- "order_id STRING," +
- "user_id STRING," +
- "order_status STRING," +
- "sku_id STRING," +
- "sku_name STRING," +
- "province_id STRING," +
- "activity_id STRING," +
- "activity_rule_id STRING," +
- "coupon_id STRING," +
- "date_id STRING," +
- "create_time STRING," +
- "operate_date_id STRING," +
- "operate_time STRING," +
- "source_id STRING," +
- "source_type_id STRING," +
- "source_type_name STRING," +
- "sku_num STRING," +
- "split_original_amount STRING," +
- "split_activity_amount STRING," +
- "split_coupon_amount STRING," +
- "split_total_amount STRING," +
- "`type` STRING," +
- "`old` map<String,String> " + MyKafkaUtil.getKafkaSinkDDL("dwd_trade_cancel_detail")
- );
-
- // TODO 5. 写出数据
- tableEnv.executeSql("INSERT INTO dwd_trade_cancel_detail SELECT * FROM filter_table");
- // TODO 6. 启动任务
- env.execute("DwdTradeCancelDetail");
- }
- }
支付成功表的主表并不是订单表或者订单明细表,而是支付表,在业务系统中存在支付表(payment_info),每行代表一个订单的支付情况(成功或失败)。
但是我们应该拿它去和下单明细表做关联,因为我们如果之后需要对支付记录做一个数据分析,就应该尽可能保证支付事实表应该有一个细的粒度和丰富的维度,而关联下单明细表之后不仅粒度商品,还增加了很多sku维度(比如商品类型、)
1)设置 ttl
支付成功事务事实表需要将业务数据库中的支付信息表 payment_info 数据与订单明细表关联。订单明细数据是在下单时生成的,经过一系列的处理进入订单明细主题,通常支付操作在下单后 15min 内完成即可,因此,支付明细数据可能比订单明细数据滞后 15min。考虑到可能的乱序问题,ttl 设置为 15min + 5s。
2)获取订单明细数据
用户必然要先下单才有可能支付成功,因此支付成功明细数据集必然是订单明细数据集的子集,所以我们直接从 dwd 层的下单明细表去读取;
3)筛选支付表数据
获取支付类型、回调时间(支付成功时间)、支付成功时间戳。
生产环境下,用户支付后,业务数据库的支付表会插入一条数据,此时的回调时间和回调内容为空。通常底层会调用第三方支付接口,接口会返回回调信息,如果支付成功则回调信息不为空,此时会更新支付表,补全回调时间和回调内容字段的值,并将 payment_status 字段的值修改为支付成功对应的状态码(本项目为 1602)。支付成功之后,支付表数据不会发生变化。因此,只要操作类型为 update 且状态码为 1602 即为支付成功数据。
由上述分析可知,支付成功对应的业务数据库变化日志应满足两个条件:
(1)payment_status 字段的值为 1602;
(2)操作类型为 update。
4)构建 MySQL-LookUp 字典表
为的是对支付方式字段做维度退化;
5)关联上述三张表形成支付成功宽表,写入 Kafka 支付成功主题
支付成功业务过程的最细粒度为一个 sku 的支付成功记录,payment_info 表的粒度与最细粒度相同,将其作为主表。
(1) payment_info 表在订单明细表中必然存在对应数据,主表不存在独有数据,因此通过内连接与订单明细表关联;
(2) 与字典表的关联是为了获取 payment_type 对应的支付类型名称,主表不存在独有数据,通过内连接与字典表关联。下文与字典表的关联同理,不再赘述。
注意:这个程序中我们需要消费三张表,其中两张来组 Kafka 需要我们指定消费者组,所以我们这里可以把它们的组设置为一样的,毕竟都在一个程序里:
- // TODO 2. 读取 topic_db 数据(这里的消费者组id尽量和下面保持一致)
- tableEnv.executeSql(MyKafkaUtil.getTopicDb("dwd_trade_pay_detail_suc"));
- // TODO 3. 过滤出支付成功的数据
- Table paymentInfo = tableEnv.sqlQuery("SELECT " +
- "data['user_id'] user_id, " +
- "data['order_id'] order_id, " +
- "data['payment_type'] payment_type, " +
- "data['callback_time'] callback_time, " +
- "pt " + // 用来和 lookup 关联
- "FROM topic_db " +
- "WHERE `database` = `gmall` " +
- "AND `table` = `payment_info` " +
- "AND `type` = 'update' " +
- "AND data['payment_status'] = '1602'"
- );
- tableEnv.createTemporaryView("payment_info",paymentInfo);
- // TODO 3. 过滤出支付成功的数据
- Table paymentInfo = tableEnv.sqlQuery("SELECT " +
- "data['user_id'] user_id, " +
- "data['order_id'] order_id, " +
- "data['payment_type'] payment_type, " +
- "data['callback_time'] callback_time, " +
- "pt " + // 用来和 lookup 关联
- "FROM topic_db " +
- "WHERE `database` = `gmall` " +
- "AND `table` = `payment_info` " +
- "AND `type` = 'update' " +
- "AND data['payment_status'] = '1602'"
- );
- tableEnv.createTemporaryView("payment_info",paymentInfo);
- // TODO 5. 读取 mysql lookup 表
- tableEnv.executeSql(MysqlUtil.getBaseDicLookUpDDL());
关联支付成功数据、下单明细和字典表:
- // TODO 6. 关联 3 张表
- Table resultTable = tableEnv.sqlQuery("SELECT " +
- " od.id order_detail_id," +
- " od.order_id," +
- " od.user_id," +
- " od.sku_id," +
- " od.sku_name," +
- " od.province_id," +
- " od.activity_id," +
- " od.activity_rule_id," +
- " od.coupon_id," +
- " pi.payment_type payment_type_code," +
- " dic.dic_name payment_type_name," +
- " pi.callback_time," +
- " od.source_id," +
- " od.source_type_code," +
- " od.source_type_name," +
- " od.sku_num," +
- " od.split_activity_amount," +
- " od.split_coupon_amount," +
- " od.split_total_amount split_payment_amount " +
- "FROM payment_info pi" +
- "join dwd_trade_order_detail od" +
- "on pi.order_id = od.order_id" +
- "join `base_dic` for system_time as of pi.proc_time as dic" +
- "on pi.payment_type = dic.dic_code ");
- tableEnv.createTemporaryView("result_table",resultTable);
- // TODO 7. 创建 kafka 支付成功表
- tableEnv.executeSql(
- "CREATE TABLE dwd_trade_pay_detail_suc (" +
- " order_detail_id string," +
- " order_id string," +
- " user_id string," +
- " sku_id string," +
- " sku_name string," +
- " province_id string," +
- " activity_id string," +
- " activity_rule_id string," +
- " coupon_id string," +
- " payment_type_code string," +
- " payment_type_name string," +
- " callback_time string," +
- " source_id string," +
- " source_type_code string," +
- " source_type_name string," +
- " sku_num string," +
- " split_activity_amount string," +
- " split_coupon_amount string," +
- " split_payment_amount string," +
- " primary key(order_detail_id) not enforced "
- + MyKafkaUtil.getUpsertKafkaDDL("dwd_trade_pay_detail_suc")
- );
这里加主键是为了方便让相同 key 到一个分区中,方便下游做去重(也可以不加);
最后的 env 不需要执行,因为我们执行的是 tableEnv ,如果执行 env 的话会有警告:没有操作可执行;
退单指的是支付成功后取消订单的操作,而取消订单指的是还没支付时取消订单的操作;
退单并不需要关联下单事务事实表,因为退单表的粒度本身就是商品(退单表包含 sku_id,有了这个 sku_id 足够之后去DIM关联 sku 维表了);我们这里需要关联订单表,为的是从中获得一些分析维度(比如 province_id,我认为还可以添加活动id和优惠券id等维度外键)。
退单操作会影响到订单表(order_status 从 1002['已支付'] 变为 1005['退款中']),所以我们可以在过滤订单表的时候,同时过滤出退单的数据,减少性能的消耗;
注意:这里我们关联的是订单表,因为退单肯定会触发修改订单表,它俩是几乎同时发生的,所以我们可以一次性从 topic_db 中读取出来。也没必要去订单预处理表去读(订单预处理表还依赖订单表,所以根本没必要)。
所以,大致思路就是:
- // TODO 2. 读取 topic_db 数据
- tableEnv.executeSql(MyKafkaUtil.getTopicDb("dwd_trade_order_refund"));
-
- // TODO 3. 过滤出退单表
- Table refundTable = tableEnv.sqlQuery("SELECT " +
- "data['d'] id, " +
- "data['user_id'] user_id, " +
- "data['order_id'] order_id, " +
- "data['sku_id'] sku_id, " +
- "data['refund_type'] refund_type, " +
- "data['refund_num'] refund_num, " +
- "data['refund_amount'] refund_amount, " +
- "data['refund_reason_type'] refund_reason_type, " +
- "data['refund_reason_txt'] refund_reason_txt, " +
- "data['create_time'] create_time, " +
- "pt " +
- "FROM topic_db " +
- "WHERE `database`='gmall' " +
- "AND `table`='order_refund_info' "
- );
- tableEnv.createTemporaryView("order_refund_info",refundTable);
- // TODO 4. 过滤出订单表中的退单数据
- Table orderInfoRefund = tableEnv.sqlQuery("SELECT " +
- "data['id'] id," +
- "data['province_id'] province_id," +
- "`old` " +
- "FROM topic_db " +
- "WHERE `database`='gmall' " +
- "AND `table`='order_info ' " +
- "AND `type` = 'update' " +
- "AND order_status = '1005' " +
- "AND old['order_status'] is not null"
- );
- tableEnv.createTemporaryView("order_info_refund",orderInfoRefund);
为的是获得退款类型和退款原因类型
- // TODO 5. 读取 mysql lookup 表
- tableEnv.executeSql(MysqlUtil.getBaseDicLookUpDDL());
- // TODO 6. 关联 3 张表
- Table resultTable = tableEnv.sqlQuery("SELECT " +
- "ri.id," +
- "ri.user_id," +
- "ri.order_id," +
- "ri.sku_id," +
- "oi.province_id," +
- "date_format(ri.create_time,'yyyy-MM-dd') date_id," +
- "ri.create_time," +
- "ri.refund_type," +
- "type_dic.dic_name," +
- "ri.refund_reason_type," +
- "reason_dic.dic_name," +
- "ri.refund_reason_txt," +
- "ri.refund_num," +
- "ri.refund_amount " +
- "from order_refund_info ri" +
- "join " +
- "order_info_refund oi" +
- "on ri.order_id = oi.id" +
- "join " +
- "base_dic for system_time as of ri.proc_time as type_dic" +
- "on ri.refund_type = type_dic.dic_code" +
- "join" +
- "base_dic for system_time as of ri.proc_time as reason_dic" +
- "on ri.refund_reason_type=reason_dic.dic_code"
- );
- tableEnv.createTemporaryView("result_table",resultTable);
- // TODO 7. 创建 Kafka 退单事务事实表
- tableEnv.executeSql("CREATE TABLE dwd_trade_order_refund (" +
- "id string," +
- "user_id string," +
- "order_id string," +
- "sku_id string," +
- "province_id string," +
- "date_id string," +
- "create_time string," +
- "refund_type_code string," +
- "refund_type_name string," +
- "refund_reason_type_code string," +
- "refund_reason_type_name string," +
- "refund_reason_txt string," +
- "refund_num string," +
- "refund_amount string," +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_trade_order_refund")
- );
-
- // TODO 8. 写入数据
- tableEnv.executeSql("INSERT INTO dwd_trade_order_refund SELECT * FROM result_table");
退款成功这个业务过程会影响到退款表(插入数据)、订单表(订单状态)和退单表(退款状态)
主要任务:
补充:
思路和前面都是一样的:
一样的套路,不解释了
- public class DwdTradeRefundPaySuc {
-
- public static void main(String[] args) {
- // TODO 1. 获取执行环境
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(1); // 生产环境中设置为kafka主题的分区数
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // 1.1 开启checkpoint
- env.enableCheckpointing(5 * 60000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointStorage("hdfs://hadoop102:8020/s/ck");
- env.getCheckpointConfig().setCheckpointTimeout(10 * 60000L);
- env.getCheckpointConfig().setMaxConcurrentCheckpoints(2); // 设置最大共存的checkpoint数量
- env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3,5000L)); // 固定频率重启: 尝试3次重启,每5s重启一次
-
- // 1.2 设置状态后端
- env.setStateBackend(new HashMapStateBackend());
-
- // 1.3 设置状态 ttl
- tableEnv.getConfig().setIdleStateRetention(Duration.ofSeconds(5));
-
- // TODO 2. 读取 topic_db 数据
- tableEnv.executeSql(MyKafkaUtil.getTopicDb("dwd_trade_refund_pay_suc"));
-
- // TODO 3. 过滤出退款表
- Table refundPayment = tableEnv.sqlQuery("SELECT " +
- "data['id'] id," +
- "data['order_id'] order_id," +
- "data['sku_id'] sku_id," +
- "data['payment_type'] payment_type," +
- "data['callback_time'] callback_time," +
- "data['total_amount'] total_amount," +
- "pt " +
- "FROM topic_db " +
- "WHERE `database`='gmall' " +
- "AND `table` = 'refund_payment' " +
- "AND data['refund_status']='0705' " +
- "AND `type` = 'update' " +
- "AND old['refund_status'] is not null"
- );
- tableEnv.createTemporaryView("refund_payment",refundPayment);
-
- // TODO 4. 过滤出订单表中退款成功的数据
- Table orderRefundSuc = tableEnv.sqlQuery("SELECT " +
- "data['id'] id, " +
- "data['user_id'] user_id, " +
- "data['province_id'] province_id, " +
- "pt " +
- "FROM topic_db " +
- "WHERE `database`='gmall' " +
- "AND `table`='order_info' " +
- "AND `type` = 'update'" +
- "AND data['order_status'] = '1006' " +
- "AND old['order_status'] is not null"
- );
- tableEnv.createTemporaryView("order_refund_suc",orderRefundSuc);
-
- // TODO 5. 过滤出退单表中退款成功的数据
- Table refundSuc = tableEnv.sqlQuery("SELECT " +
- "data['order_id'] order_id, " +
- "data['sku_id'] sku_id, " +
- "data['refund_num'] refund_num, " +
- "pt " +
- "FROM topic_db " +
- "WHERE `database`='gmall' " +
- "AND `table`='order_refund_info' " +
- "AND `type` = 'update'" +
- "AND data['order_status'] = '0705' " +
- "AND old['refund_status'] is not null"
- );
- tableEnv.createTemporaryView("order_refund_info",refundSuc);
-
- // TODO 6. 创建 MySQL lookup表
- tableEnv.executeSql(MysqlUtil.getBaseDicLookUpDDL());
-
-
- // TODO 7. 关联 4 张表
- Table resultTable = tableEnv.sqlQuery("select" +
- "rp.id," +
- "oi.user_id," +
- "rp.order_id," +
- "rp.sku_id," +
- "oi.province_id," +
- "rp.payment_type," +
- "dic.dic_name payment_type_name," +
- "date_format(rp.callback_time,'yyyy-MM-dd') date_id," +
- "rp.callback_time," +
- "ri.refund_num," +
- "rp.total_amount," +
- "from refund_payment rp " +
- "join " +
- "order_info oi" +
- "on rp.order_id = oi.id" +
- "join" +
- "order_refund_info ri" +
- "on rp.order_id = ri.order_id" +
- "and rp.sku_id = ri.sku_id" +
- "join " +
- "base_dic for system_time as of rp.proc_time as dic" +
- "on rp.payment_type = dic.dic_code");
- tableEnv.createTemporaryView("result_table", resultTable);
-
- // TODO 8. 创建 Kafka 退款成功事务事实表
- tableEnv.executeSql("create table dwd_trade_refund_pay_suc(" +
- "id string," +
- "user_id string," +
- "order_id string," +
- "sku_id string," +
- "province_id string," +
- "payment_type_code string," +
- "payment_type_name string," +
- "date_id string," +
- "callback_time string," +
- "refund_num string," +
- "refund_amount string "+
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_trade_refund_pay_suc"));
-
- // TODO 8. 写入数据
- tableEnv.executeSql("INSERT INTO dwd_trade_refund_pay_suc SELECT * FROM result_table");
- }
- }
我们先看看业务系统中 coupon_use 中的数据:
这张表的粒度是一个订单,每行数据代表的哪个用户在什么时候领取的券,什么时候使用了券( using_time 表示使用时间,used_time 表示支付时间),什么时候过期等信息;
每发生优惠券领取这样一个业务过程,coupon_use 数据库中就会 insert 一条新的数据,所以我们只需要直接从 topic_db 中过滤出 type = 'insert' 的数据即可;
而关于优惠券下单使用(下单)和优惠券使用(支付)这两个业务过程除了过滤出 type = 'update' 之外,只需要分别对 using_time(下单) 和 used_time(支付) 是否为 null 进行过滤即可;
总结:
很简单,不多废话了
- public class DwdToolCouponGet {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // TODO 2. 状态后端设置
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.setRestartStrategy(RestartStrategies.failureRateRestart(
- 3, Time.days(1), Time.minutes(1)
- ));
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage(
- "hdfs://hadoop102:8020/ck"
- );
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table `topic_db`(" +
- "`database` string," +
- "`table` string," +
- "`data` map<string, string>," +
- "`type` string," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_tool_coupon_get"));
-
- // TODO 4. 读取优惠券领用数据,封装为表
- Table resultTable = tableEnv.sqlQuery("select" +
- "data['id']," +
- "data['coupon_id']," +
- "data['user_id']," +
- "date_format(data['get_time'],'yyyy-MM-dd') date_id," +
- "data['get_time']," +
- "ts" +
- "from topic_db" +
- "where `table` = 'coupon_use'" +
- "and `type` = 'insert'");
- tableEnv.createTemporaryView("result_table", resultTable);
-
- // TODO 5. 建立 Kafka-Connector dwd_tool_coupon_get 表
- tableEnv.executeSql("create table dwd_tool_coupon_get (" +
- "id string," +
- "coupon_id string," +
- "user_id string," +
- "date_id string," +
- "get_time string," +
- "ts string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_tool_coupon_get"));
-
- // TODO 6. 将数据写入 Kafka-Connector 表
- tableEnv.executeSql("insert into dwd_tool_coupon_get select * from result_table");
- }
- }
- public class DwdToolCouponOrder {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // TODO 2. 状态后端设置
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.setRestartStrategy(RestartStrategies.failureRateRestart(
- 3, Time.days(1), Time.minutes(1)
- ));
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage(
- "hdfs://hadoop102:8020/ck"
- );
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table `topic_db` (" +
- "`database` string," +
- "`table` string," +
- "`data` map<string, string>," +
- "`type` string," +
- "`old` map<string, string>," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_tool_coupon_order"));
-
- // TODO 4. 读取优惠券领用表数据,筛选满足条件的优惠券下单数据
- Table couponUseOrder = tableEnv.sqlQuery("select" +
- "data['id'] id," +
- "data['coupon_id'] coupon_id," +
- "data['user_id'] user_id," +
- "data['order_id'] order_id," +
- "date_format(data['using_time'],'yyyy-MM-dd') date_id," +
- "data['using_time'] using_time," +
- "ts" +
- "from topic_db" +
- "where `table` = 'coupon_use'" +
- "and `type` = 'update'" +
- "and data['coupon_status'] = '1402'" +
- "and `old`['coupon_status'] = '1401'");
-
- tableEnv.createTemporaryView("result_table", couponUseOrder);
-
- // TODO 5. 建立 Kafka-Connector dwd_tool_coupon_order 表
- tableEnv.executeSql("create table dwd_tool_coupon_order(" +
- "id string," +
- "coupon_id string," +
- "user_id string," +
- "order_id string," +
- "date_id string," +
- "order_time string," +
- "ts string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_tool_coupon_order"));
-
- // TODO 6. 将数据写入 Kafka-Connector 表
- tableEnv.executeSql("" +
- "insert into dwd_tool_coupon_order select " +
- "id," +
- "coupon_id," +
- "user_id," +
- "order_id," +
- "date_id," +
- "using_time order_time," +
- "ts from result_table");
- }
- }
- public class DwdToolCouponPay {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // TODO 2. 状态后端设置
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.setRestartStrategy(RestartStrategies.failureRateRestart(
- 3, Time.days(1), Time.minutes(1)
- ));
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage(
- "hdfs://hadoop102:8020/ck"
- );
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table `topic_db` (" +
- "`database` string," +
- "`table` string," +
- "`data` map<string, string>," +
- "`type` string," +
- "`old` string," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_tool_coupon_pay"));
-
- // TODO 4. 读取优惠券领用表数据,筛选优惠券使用(支付)数据
- Table couponUsePay = tableEnv.sqlQuery("select" +
- "data['id'] id," +
- "data['coupon_id'] coupon_id," +
- "data['user_id'] user_id," +
- "data['order_id'] order_id," +
- "date_format(data['used_time'],'yyyy-MM-dd') date_id," +
- "data['used_time'] used_time," +
- "`old`," +
- "ts" +
- "from topic_db" +
- "where `table` = 'coupon_use'" +
- "and `type` = 'update'" +
- "and data['used_time'] is not null");
-
- tableEnv.createTemporaryView("coupon_use_pay", couponUsePay);
-
- // TODO 5. 建立 Kafka-Connector dwd_tool_coupon_order 表
- tableEnv.executeSql("create table dwd_tool_coupon_pay(" +
- "id string," +
- "coupon_id string," +
- "user_id string," +
- "order_id string," +
- "date_id string," +
- "payment_time string," +
- "ts string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_tool_coupon_pay"));
-
- // TODO 6. 将数据写入 Kafka-Connector 表
- tableEnv.executeSql("" +
- "insert into dwd_tool_coupon_pay select " +
- "id," +
- "coupon_id," +
- "user_id," +
- "order_id," +
- "date_id," +
- "used_time payment_time," +
- "ts from coupon_use_pay");
- }
- }
我可以看一下业务系统中的收藏表:
可以看到,收藏表的粒度是商品,每收藏一件商品 favor_info 表就会 insert 一条数据,每取消收藏一件商品,并不会删除记录,而是将收藏表的 is_cancel 字段设为 1 并给 cancel_time 字段补充值;所以过滤条件很简单:
- public class DwdInteractionFavorAdd {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // TODO 2. 状态后端设置
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.setRestartStrategy(RestartStrategies.failureRateRestart(
- 3, Time.days(1), Time.minutes(1)
- ));
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage(
- "hdfs://hadoop102:8020/ck"
- );
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table topic_db(" +
- "`database` string," +
- "`table` string," +
- "`type` string," +
- "`data` map<string, string>," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_interaction_favor_add"));
-
- // TODO 4. 读取收藏表数据
- Table favorInfo = tableEnv.sqlQuery("select" +
- "data['id'] id," +
- "data['user_id'] user_id," +
- "data['sku_id'] sku_id," +
- "date_format(data['create_time'],'yyyy-MM-dd') date_id," +
- "data['create_time'] create_time," +
- "ts" +
- "from topic_db" +
- "where `table` = 'favor_info'" +
- "and (`type` = 'insert' or (`type` = 'insert' and data['is_cancel'] = '1'))");
- tableEnv.createTemporaryView("favor_info", favorInfo);
-
- // TODO 5. 创建 Kafka-Connector dwd_interaction_favor_add 表
- tableEnv.executeSql("create table dwd_interaction_favor_add (" +
- "id string," +
- "user_id string," +
- "sku_id string," +
- "date_id string," +
- "create_time string," +
- "ts string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_interaction_favor_add"));
-
- // TODO 6. 将数据写入 Kafka-Connector 表
- tableEnv.executeSql("" +
- "insert into dwd_interaction_favor_add select * from favor_info");
- }
- }
任务:建立 MySQL-Lookup 字典表,读取评论表数据,关联字典表以获取评价(好评、中评、差评、自动),将结果写入 Kafka 评价主题
所以这张表的逻辑页很简单:
- public class DwdInteractionComment {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
-
- // 获取配置对象
- Configuration configuration = tableEnv.getConfig().getConfiguration();
- // 为表关联时状态中存储的数据设置过期时间
- configuration.setString("table.exec.state.ttl", "5 s");
-
- // TODO 2. 状态后端设置
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.setRestartStrategy(RestartStrategies.failureRateRestart(
- 3, Time.days(1), Time.minutes(1)
- ));
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage(
- "hdfs://hadoop102:8020/ck"
- );
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table topic_db(" +
- "`database` string," +
- "`table` string," +
- "`type` string," +
- "`data` map<string, string>," +
- "`proc_time` as PROCTIME()," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_interaction_comment"));
-
- // TODO 4. 读取评论表数据
- Table commentInfo = tableEnv.sqlQuery("select" +
- "data['id'] id," +
- "data['user_id'] user_id," +
- "data['sku_id'] sku_id," +
- "data['order_id'] order_id," +
- "data['create_time'] create_time," +
- "data['appraise'] appraise," +
- "proc_time," +
- "ts" +
- "from topic_db" +
- "where `table` = 'comment_info'" +
- "and `type` = 'insert'");
- tableEnv.createTemporaryView("comment_info", commentInfo);
-
- // TODO 5. 建立 MySQL-LookUp 字典表
- tableEnv.executeSql(MysqlUtil.getBaseDicLookUpDDL());
-
- // TODO 6. 关联两张表
- Table resultTable = tableEnv.sqlQuery("select" +
- "ci.id," +
- "ci.user_id," +
- "ci.sku_id," +
- "ci.order_id," +
- "date_format(ci.create_time,'yyyy-MM-dd') date_id," +
- "ci.create_time," +
- "ci.appraise," +
- "dic.dic_name," +
- "ts" +
- "from comment_info ci" +
- "join" +
- "base_dic for system_time as of ci.proc_time as dic" +
- "on ci.appraise = dic.dic_code");
- tableEnv.createTemporaryView("result_table", resultTable);
-
- // TODO 7. 建立 Kafka-Connector dwd_interaction_comment 表
- tableEnv.executeSql("create table dwd_interaction_comment(" +
- "id string," +
- "user_id string," +
- "sku_id string," +
- "order_id string," +
- "date_id string," +
- "create_time string," +
- "appraise_code string," +
- "appraise_name string," +
- "ts string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_interaction_comment"));
-
- // TODO 8. 将关联结果写入 Kafka-Connector 表
- tableEnv.executeSql("" +
- "insert into dwd_interaction_comment select * from result_table");
- }
- }
主要任务:读取用户表数据,获取注册时间,将用户注册信息写入 Kafka 用户注册主题
如何分辨是新用户还是老用户很简单:新 insert 进来的都是新用户,所以思路很简单:
- public class DwdUserRegister {
- public static void main(String[] args) throws Exception {
-
- // TODO 1. 环境准备
- StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
- env.setParallelism(4);
- StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
- tableEnv.getConfig().setLocalTimeZone(ZoneId.of("GMT+8"));
-
- // TODO 2. 启用状态后端
- env.enableCheckpointing(3000L, CheckpointingMode.EXACTLY_ONCE);
- env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L);
- env.getCheckpointConfig().enableExternalizedCheckpoints(
- CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
- );
- env.getCheckpointConfig().setMinPauseBetweenCheckpoints(3000L);
- env.setRestartStrategy(
- RestartStrategies.failureRateRestart(3, Time.days(1L), Time.minutes(3L))
- );
- env.setStateBackend(new HashMapStateBackend());
- env.getCheckpointConfig().setCheckpointStorage("hdfs://hadoop102:8020/ck");
- System.setProperty("HADOOP_USER_NAME", "lyh");
-
- // TODO 3. 从 Kafka 读取业务数据,封装为 Flink SQL 表
- tableEnv.executeSql("create table topic_db(" +
- "`database` string," +
- "`table` string," +
- "`type` string," +
- "`data` map<string, string>," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaDDL("topic_db", "dwd_trade_order_detail"));
-
- // TODO 4. 读取用户表数据
- Table userInfo = tableEnv.sqlQuery("select" +
- "data['id'] user_id," +
- "data['create_time'] create_time," +
- "ts" +
- "from topic_db" +
- "where `table` = 'user_info'" +
- "and `type` = 'insert'");
- tableEnv.createTemporaryView("user_info", userInfo);
-
- // TODO 5. 创建 Kafka-Connector dwd_user_register 表
- tableEnv.executeSql("create table `dwd_user_register`(" +
- "`user_id` string," +
- "`date_id` string," +
- "`create_time` string," +
- "`ts` string" +
- ")" + MyKafkaUtil.getKafkaSinkDDL("dwd_user_register"));
-
- // TODO 6. 将输入写入 Kafka-Connector 表
- tableEnv.executeSql("insert into dwd_user_register" +
- "select " +
- "user_id," +
- "date_format(create_time, 'yyyy-MM-dd') date_id," +
- "create_time," +
- "ts" +
- "from user_info");
-
- }
- }
至此,DWD 层搭建完毕,实时数仓的 DWD 层没有离线数仓那么多分类(周期快照、累积快照等),所以做起来虽然不是那么简单,但是有迹可循;
明天开始就是 DWS 层的搭建了,终于可以用到 clickhouse 了,浅浅期待一下~
这次做实时数仓项目比离线快很多,但是并不是囫囵吞枣,看视频三年了,是在水视频过任务还是自己真心想理解一个项目自己很清楚。这个实时数仓和之前学的离线数仓的数据源都是一样的,对这个模拟出来的业务系统比较熟悉了,而且离线和实时有很多共同点所以学的很快,但是毕竟时间也花费了不少,每天的早八晚八;
很多人做项目也好,学新技术也好,看到几十个小时的时长就想着过任务一样,自己骗自己看过视频就算学会了,就算过去了,好像再也不用看了一样。前期基础不好好学,后期很难有那种"灵光乍现"、"触类旁通"的感觉,因为知识储备就不够。就像我自己每天说要看书(今天下单一本《乡土中国》),但是老想想着刷刷抖音,骗自己抖音的地摊文化更有营养,到头来现在在写作或者表达的时候还是经常词穷;
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。