· 10 years ago · Aug 22, 2016, 11:36 PM
1""""
2@desc Information about location-driven radio wakeups across demographic cuts.
3"""
4from __future__ import absolute_import
5from __future__ import division
6from __future__ import print_function
7
8from dataswarm.operators import (
9 GlobalDefaults,
10 HiveQLOperator,
11 WaitForHiveOperator,
12 HivePublishSignalTableOperator,
13)
14
15GlobalDefaults.set(
16 user='sbis',
17 schedule='@daily',
18 pool='locations.highpri_pipelines',
19 fail_on_future=False,
20 depends_on_past=False,
21 partition='ds=<DATEID>',
22 dep_list=[],
23 max_concurrent=5,
24)
25
26
27populate_table_aggregated_location_coverage_cycles = {}
28
29###############################################################################
30# PIPELINE START #
31###############################################################################
32create_mobile_location_driven_radio_wakeups = HiveQLOperator(
33 dep_list=[],
34 hive_query="""
35 CREATE TABLE IF NOT EXISTS <TABLE:mobile_location_driven_radio_wakeups> (
36
37 user_id BIGINT,
38
39 is_employee TINYINT,
40 mobile_l1 TINYINT,
41 mobile_l7 TINYINT,
42 mobile_l28 TINYINT,
43 country STRING,
44
45 app_id STRING,
46 device_id STRING,
47 app_version STRING,
48 carrier STRING,
49 is_background SMALLINT,
50 call_path STRING,
51 connection_quality_histogram MAP<STRING, BIGINT>,
52 connection_type_histogram MAP<STRING, BIGINT>,
53 connection_subtype_histogram MAP<STRING, BIGINT>,
54 session_count BIGINT,
55 total_period_duration BIGINT,
56 wakeup_count INT,
57 norm_wakeup_count INT,
58 request_count INT,
59 norm_request_count INT
60 )
61
62 PARTITIONED BY (
63 ds STRING,
64 interface STRING)
65 TBLPROPERTIES('RETENTION'='30');
66 """,
67)
68
69wait_for_dim_all_users_ppt = WaitForHiveOperator(
70 table='dim_all_users_ppt_signal:bi',
71)
72
73wait_for_mobile_power_metrics_liger_attribution_persist = WaitForHiveOperator(
74 table='mobile_power_metrics_liger_attribution_persist_signal:mobile',
75)
76
77populate_table_aggregated_location_coverage_cycles['android'] = HiveQLOperator(
78 dep_list=[wait_for_dim_all_users_ppt,
79 wait_for_mobile_power_metrics_liger_attribution_persist,
80 create_mobile_location_driven_radio_wakeups],
81 hive_query="""
82 SET mapred.min.split.size=64000000;
83 SET mapred.max.split.size=64000000;
84 SET hive.exec.reducers.bytes.per.reducer=125000000;
85 SET hive.exec.reducers.max = 512;
86
87 INSERT OVERWRITE TABLE <TABLE:mobile_location_driven_radio_wakeups>
88 PARTITION (
89 ds = '<DATEID>',
90 interface = 'android'
91 )
92 SELECT
93 mpm_fb4a.user_id,
94
95 dau_ppt.is_employee,
96 dau_ppt.mobile_l1,
97 dau_ppt.mobile_l7,
98 dau_ppt.mobile_l28,
99 dau_ppt.country,
100
101 mpm_fb4a.app_id,
102 mpm_fb4a.device_id,
103 mpm_fb4a.app_version,
104 mpm_fb4a.carrier,
105 mpm_fb4a.is_background,
106 mpm_fb4a.call_path,
107 mpm_fb4a.connection_quality_histogram,
108 mpm_fb4a.connection_type_histogram,
109 mpm_fb4a.connection_subtype_histogram,
110 mpm_fb4a.session_count,
111 mpm_fb4a.total_period_duration,
112 mpm_fb4a.wakeup_count,
113 mpm_fb4a.norm_wakeup_count,
114 mpm_fb4a.request_count,
115 mpm_fb4a.norm_request_count
116 FROM (
117 SELECT
118 user_id,
119 app_id,
120 device_id,
121 app_version,
122 carrier,
123 is_background,
124 call_path,
125 FB_HISTOGRAM(connection_quality) AS connection_quality_histogram,
126 FB_HISTOGRAM(connection_type) AS connection_type_histogram,
127 FB_HISTOGRAM(connection_subtype) AS connection_subtype_histogram,
128 SUM(session_count) AS session_count,
129 SUM(period_duration) AS total_period_duration, -- overlaps?
130 SUM(wakeup_count) AS wakeup_count,
131 SUM(norm_wakeup_count) AS norm_wakeup_count,
132 SUM(request_count) AS request_count,
133 SUM(norm_request_count) AS norm_request_count
134 FROM
135 mobile_power_metrics_liger_attribution_persist:mobile
136 WHERE
137 ds = '<DATEID>'
138 AND app_id = '350685531728'
139 AND (call_path LIKE '%location%' OR call_path LIKE '%Location%')
140 GROUP BY
141 1, 2, 3, 4, 5, 6, 7
142 ) mpm_fb4a
143 JOIN
144 dim_all_users_ppt:bi dau_ppt
145 ON
146 mpm_fb4a.user_id = dau_ppt.userid
147 AND dau_ppt.ds = '<DATEID>'
148 """
149)
150
151job_end = HivePublishSignalTableOperator(
152 base_table_name='mobile_location_driven_radio_wakeups',
153 dep_list=[
154 create_mobile_location_driven_radio_wakeups,
155 ] + populate_table_aggregated_location_coverage_cycles.values(),
156)