· 9 years ago · Oct 17, 2016, 01:34 PM
1from __future__ import absolute_import
2from __future__ import division
3from __future__ import print_function
4from dataswarm.operators import (
5 WaitForHiveOperator,
6 HiveQLOperator,
7 GlobalDefaults
8)
9
10## set the global defaults;
11## put as many parameters here that will be used by all the tasks
12GlobalDefaults.set(
13 user='willnuland',
14 schedule='@daily',
15 pool='di.lowpri_pipelines',
16)
17
18dim_all_users = WaitForHiveOperator(
19 table='dim_all_users:bi',
20 ds='<DATEID>'
21)
22
23dyi_generation = WaitForHiveOperator(
24 table='dyi_generation:si',
25 ds='<DATEID>'
26)
27
28create_hive_table = HiveQLOperator(
29 dep_list=[dim_all_users, dyi_generation],
30 hive_query=""
31 SELECT
32 dau.country COUNT (1) as country
33FROM
34 dyi_generation dyi
35JOIN dim_all_users:bi dau
36ON dyi.account_uid = dau.userid
37WHERE
38 dyi.ds = '<DATEID>'
39 and dau.ds = '<DATEID>'
40 and dyi.event = 'completed'
41 GROUP by dau.country""
42 CREATE TABLE IF NOT EXISTS <TABLE:dyi_by_country> (
43 cnt BIGINT
44 ) PARTITIONED BY (ds STRING)
45 TBLPROPERTIES('RETENTION'='90');
46
47 INSERT OVERWRITE TABLE <TABLE:dyi_by_country>
48 PARTITION (ds='<DATEID>')
49
50 """
51)