Skip to content

upload

upload

Dataset upload and data operations.

index_dataset async

index_dataset(
    auth: DomoAuth,
    dataset_id: str,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> ResponseGetData

manually index a dataset

Source code in src/crew_dcs/routes/dataset/upload.py
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def index_dataset(
    auth: DomoAuth,
    dataset_id: str,
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """manually index a dataset"""

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/indexes"

    body = {"dataIds": []}

    res = await gd.get_data(
        auth=auth,
        method="POST",
        body=body,
        url=url,
        context=context,
    )

    if not res.is_success:
        raise Dataset_CRUD_Error(dataset_id=dataset_id, res=res)

    return res

index_status async

index_status(
    auth: DomoAuth,
    dataset_id: str,
    index_id: str,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> ResponseGetData

get the completion status of an index

Source code in src/crew_dcs/routes/dataset/upload.py
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def index_status(
    auth: DomoAuth,
    dataset_id: str,
    index_id: str,
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """get the completion status of an index"""

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/indexes/{index_id}/statuses"

    res = await gd.get_data(
        auth=auth,
        method="GET",
        url=url,
        context=context,
    )

    if not res.is_success:
        raise Dataset_GET_Error(dataset_id=dataset_id, res=res)

    return res

list_partitions async

list_partitions(
    auth: DomoAuth,
    dataset_id: str,
    body: dict | None = None,
    debug_loop: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs
)

List all partitions for a dataset.

Source code in src/crew_dcs/routes/dataset/upload.py
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def list_partitions(
    auth: DomoAuth,
    dataset_id: str,
    body: dict | None = None,
    debug_loop: bool = False,
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
):
    """List all partitions for a dataset."""

    body = body or generate_list_partitions_body()

    url = f"https://{auth.domo_instance}.domo.com/api/query/v1/datasources/{dataset_id}/partition/list"

    offset_params = {
        "offset": "offset",
        "limit": "limit",
    }

    def arr_fn(res) -> list[dict]:
        return res.response

    res = await gd.looper(
        auth=auth,
        method="POST",
        url=url,
        arr_fn=arr_fn,
        body=body,
        offset_params_in_body=True,
        offset_params=offset_params,
        loop_until_end=True,
        debug_loop=debug_loop,
        context=context,
    )

    if res.status == 404 and res.response == "Not Found":
        raise DatasetNotFoundError(dataset_id=dataset_id, res=res)

    if not res.is_success:
        raise Dataset_GET_Error(dataset_id=dataset_id, res=res)

    return res

upload_dataset_stage_1 async

upload_dataset_stage_1(
    auth: DomoAuth,
    dataset_id: str,
    partition_tag: str | None = None,
    *,
    context: RouteContext | None = None,
    return_raw: bool = False,
    **context_kwargs
) -> ResponseGetData

preps dataset for upload by creating an upload_id (upload session key) pass to stage 2 as a parameter

Source code in src/crew_dcs/routes/dataset/upload.py
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
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def upload_dataset_stage_1(
    auth: DomoAuth,
    dataset_id: str,
    partition_tag: str | None = None,  # synonymous with data_tag
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    return_raw: bool = False,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """preps dataset for upload by creating an upload_id (upload session key) pass to stage 2 as a parameter"""

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/uploads"

    # base body assumes no paritioning
    body = {"action": None, "appendId": None}

    params = None

    if partition_tag:
        params = {"dataTag": partition_tag}
        body.update({"appendId": "latest"})  # type: ignore

    res = await gd.get_data(
        auth=auth,
        url=url,
        method="POST",
        body=body,
        params=params,
        context=context,
    )

    if not res.is_success:
        raise UploadDataError(stage_num=1, dataset_id=dataset_id, res=res)

    if return_raw:
        return res

    upload_id = res.response.get("uploadId")

    if not upload_id:
        raise UploadDataError(
            stage_num=1,
            dataset_id=dataset_id,
            res=res,
            message="no upload_id",
        )

    res.response = upload_id

    return res

upload_dataset_stage_2_df async

upload_dataset_stage_2_df(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,
    upload_df: DataFrame,
    part_id: int = 2,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> ResponseGetData

Upload pandas DataFrame to dataset (stage 2 of upload process).

Source code in src/crew_dcs/routes/dataset/upload.py
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
166
167
168
169
170
171
172
173
174
175
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def upload_dataset_stage_2_df(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,  # must originate from  a stage_1 upload response
    upload_df: pd.DataFrame,
    part_id: int = 2,  # only necessary if streaming multiple files into the same partition (multi-part upload)
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """Upload pandas DataFrame to dataset (stage 2 of upload process)."""

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/uploads/{upload_id}/parts/{part_id}"

    body = upload_df.to_csv(header=False, index=False)

    # if debug:

    res = await gd.get_data(
        url=url,
        method="PUT",
        auth=auth,
        content_type="text/csv",
        body=body,
        context=context,
    )

    if not res.is_success:
        raise UploadDataError(stage_num=2, dataset_id=dataset_id, res=res)

    res.upload_id = upload_id
    res.dataset_id = dataset_id
    res.part_id = part_id

    return res

upload_dataset_stage_2_file async

upload_dataset_stage_2_file(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,
    data_file: TextIOWrapper | None = None,
    part_id: int = 2,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> ResponseGetData

Upload data file to dataset (stage 2 of upload process).

Source code in src/crew_dcs/routes/dataset/upload.py
 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
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def upload_dataset_stage_2_file(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,  # must originate from  a stage_1 upload response
    data_file: io.TextIOWrapper | None = None,
    # only necessary if streaming multiple files into the same partition (multi-part upload)
    part_id: int = 2,
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """Upload data file to dataset (stage 2 of upload process)."""

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/uploads/{upload_id}/parts/{part_id}"

    body = data_file

    res = await gd.get_data(
        url=url,
        method="PUT",
        auth=auth,
        content_type="text/csv",
        body=body,
        context=context,
    )

    if not res.is_success:
        raise UploadDataError(stage_num=2, dataset_id=dataset_id, res=res)

    res.upload_id = upload_id
    res.dataset_id = dataset_id
    res.part_id = part_id

    return res

upload_dataset_stage_3 async

upload_dataset_stage_3(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,
    update_method: str = "REPLACE",
    partition_tag: str | None = None,
    is_index: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> ResponseGetData

commit will close the upload session, upload_id. this request defines how the data will be loaded into Adrenaline, update_method has optional flag for indexing dataset.

Source code in src/crew_dcs/routes/dataset/upload.py
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
@gd.route_function
@log_call(
    level_name="route",
    config=LogDecoratorConfig(
        entity_extractor=DomoEntityExtractor(),
        result_processor=DomoEntityResultProcessor(),
    ),
)
async def upload_dataset_stage_3(
    auth: DomoAuth,
    dataset_id: str,
    upload_id: str,  # must originate from  a stage_1 upload response
    update_method: str = "REPLACE",  # accepts REPLACE or APPEND
    partition_tag: str | None = None,  # synonymous with data_tag
    is_index: bool = False,  # index after uploading
    *,  # Make following params keyword-only
    context: RouteContext | None = None,
    **context_kwargs,
) -> rgd.ResponseGetData:
    """commit will close the upload session, upload_id.  this request defines how the data will be loaded into Adrenaline, update_method
    has optional flag for indexing dataset.
    """
    context = RouteContext.build_context(context=context, **context_kwargs)

    url = f"https://{auth.domo_instance}.domo.com/api/data/v3/datasources/{dataset_id}/uploads/{upload_id}/commit"

    body = {"index": is_index, "action": update_method}

    if partition_tag:
        body.update(
            {
                "action": "APPEND",
                "dataTag": partition_tag,
                "appendId": "latest" if partition_tag else None,
                "index": is_index,
            }
        )

    res = await gd.get_data(
        auth=auth,
        method="PUT",
        url=url,
        body=body,
        context=context,
    )

    if not res.is_success:
        raise UploadDataError(stage_num=3, dataset_id=dataset_id, res=res)

    res.upload_id = upload_id
    res.dataset_id = dataset_id

    return res