-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathstripe_analytics_pipeline.py
165 lines (148 loc) · 5.41 KB
/
stripe_analytics_pipeline.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
from typing import Optional, Tuple
import dlt
from dlt_plus.dbt_generator.utils import table_reference_adapter
from pendulum import DateTime
from stripe_analytics import (
ENDPOINTS,
INCREMENTAL_ENDPOINTS,
incremental_stripe_source,
stripe_source,
)
def load_data(
endpoints: Tuple[str, ...] = ENDPOINTS + INCREMENTAL_ENDPOINTS,
start_date: Optional[DateTime] = None,
end_date: Optional[DateTime] = None,
) -> None:
"""
This demo script uses the resources with non-incremental
loading based on "replace" mode to load all data from provided endpoints.
Args:
endpoints: A tuple of endpoint names to retrieve data from. Defaults to most popular Stripe API endpoints.
start_date: An optional start date to limit the data retrieved. Defaults to None.
end_date: An optional end date to limit the data retrieved. Defaults to None.
"""
pipeline = dlt.pipeline(
pipeline_name="stripe_analytics",
destination='bigquery',
dataset_name="stripe_pipeline",
)
source = stripe_source(
endpoints=endpoints, start_date=start_date, end_date=end_date
)
load_info = pipeline.run(source)
table_reference_adapter(
source,
"subscription",
references=[
{
"referenced_table": "event__data__object__discounts",
"columns": ["discount__id"],
"referenced_columns": ["id"],
}
],
)
table_reference_adapter(
source,
"subscription",
references=[
{
"referenced_table": "coupon",
"columns": ["discount__coupon__id"],
"referenced_columns": ["id"],
}
],
)
table_reference_adapter(
source,
"subscription",
references=[
{
"referenced_table": "price",
"columns": ["plan__id"],
"referenced_columns": ["id"],
}
],
)
table_reference_adapter(
source,
"invoice",
references=[
{
"referenced_table": "coupon",
"columns": ["discount__coupon__id"],
"referenced_columns": ["id"],
}
],
)
table_reference_adapter(
source,
"invoice",
references=[
{
"referenced_table": "event__data__object__discounts",
"columns": ["discount__id"],
"referenced_columns": ["id"],
}
],
)
print(load_info)
def load_incremental_endpoints(
endpoints: Tuple[str, ...] = INCREMENTAL_ENDPOINTS,
initial_start_date: Optional[DateTime] = None,
end_date: Optional[DateTime] = None,
) -> None:
"""
This demo script demonstrates the use of resources with incremental loading, based on the "append" mode.
This approach enables us to load all the data
for the first time and only retrieve the newest data later,
without duplicating and downloading a massive amount of data.
Make sure you're loading objects that don't change over time.
Args:
endpoints: A tuple of incremental endpoint names to retrieve data from.
Defaults to Stripe API endpoints with uneditable data.
initial_start_date: An optional parameter that specifies the initial value for dlt.sources.incremental.
If parameter is not None, then load only data that were created after initial_start_date on the first run.
Defaults to None. Format: datetime(YYYY, MM, DD).
end_date: An optional end date to limit the data retrieved.
Defaults to None. Format: datetime(YYYY, MM, DD).
"""
pipeline = dlt.pipeline(
pipeline_name="stripe_analytics",
destination='bigquery',
dataset_name="stripe_incremental",
)
# load all data on the first run that created before end_date
source = incremental_stripe_source(
endpoints=endpoints,
initial_start_date=initial_start_date,
end_date=end_date,
)
load_info = pipeline.run(source)
print(load_info)
# # load nothing, because incremental loading and end date limit
# source = incremental_stripe_source(
# endpoints=endpoints,
# initial_start_date=initial_start_date,
# end_date=end_date,
# )
# load_info = pipeline.run(source)
# print(load_info)
#
# # load only the new data that created after end_date
# source = incremental_stripe_source(
# endpoints=endpoints,
# initial_start_date=initial_start_date,
# )
# load_info = pipeline.run(source)
# print(load_info)
if __name__ == "__main__":
load_data()
# # load only data that was created during the period between the Jan 1, 2024 (incl.), and the Feb 1, 2024 (not incl.).
# from pendulum import datetime
# load_data(start_date=datetime(2024, 1, 1), end_date=datetime(2024, 2, 1))
# # load only data that was created during the period between the May 3, 2023 (incl.), and the March 1, 2024 (not incl.).
# load_incremental_endpoints(
# endpoints=("Event",),
# initial_start_date=datetime(2023, 5, 3),
# end_date=datetime(2024, 3, 1),
# )