Skip to content

stream

stream

DomoStream dataclass

DomoStream(
    auth: DomoAuth,
    id: str,
    raw: dict,
    parent: Any,
    transport_description: str = None,
    transport_version: int = None,
    update_method: str = None,
    data_provider_name: str = None,
    data_provider_key: str = None,
    account_id: str = None,
    account_display_name: str = None,
    account_userid: str = None,
    configuration_tables: list[str] = list(),
    configuration_query: str | None = None,
    _typed_config: StreamConfig_Base | None = None,
    Schedule: DomoSchedule = None,
    Account: Any = None,
)

Bases: DomoEntity

A class for interacting with a Domo Stream (dataset connector)

bucket property

bucket: str | None

Get S3 bucket or cloud storage location from stream configuration.

Works for AWS S3 connectors.

Returns:

Type Description
str | None

Bucket name/path or None if not available

Example

stream.bucket "my-s3-bucket"

database property

database: str | None

Get database name from stream configuration (cross-platform).

Works across database providers: - Snowflake (all variants) - AWS Athena - PostgreSQL

Returns:

Type Description
str | None

Database name or None if not available

Example

stream.database "SA_PRD"

dataset_id property

dataset_id: str

Get dataset ID for this stream.

For Domo-to-Domo dataset copy connectors, returns the source dataset ID. For all other streams, returns the parent dataset ID.

Returns:

Type Description
str

Dataset ID string

Example

stream.dataset_id "abc123"

display_url property

display_url

Generate URL to view this stream in the Domo UI

file_url property

file_url: str | None

Get file URL from stream configuration.

Works for file-based connectors like Domo CSV.

Returns:

Type Description
str | None

File URL or None if not available

Example

stream.file_url "https://example.com/data.csv"

host property

host: str | None

Get host/server address from stream configuration.

Works for database and API connectors: - PostgreSQL (database host) - SharePoint Online (site URL)

Returns:

Type Description
str | None

Host address/URL or None if not available

Example

stream.host "mydb.example.com"

port property

port: str | None

Get port number from stream configuration.

Works for database connectors like PostgreSQL.

Returns:

Type Description
str | None

Port number as string or None if not available

Example

stream.port "5432"

report_id property

report_id: str | None

Get report identifier from stream configuration (cross-platform).

Works for reporting/analytics platforms: - Adobe Analytics (report suite ID) - Qualtrics (survey ID)

Returns:

Type Description
str | None

Report/survey identifier or None if not available

Example

stream.report_id "report_suite_12345"

schema property

schema: str | None

Get schema name from stream configuration (cross-platform).

Works for providers that support schemas: - Snowflake (with keypair auth) - PostgreSQL

Returns:

Type Description
str | None

Schema name or None if not available

Example

stream.schema "PUBLIC"

source_account property

source_account: str | None

Get the source system account identifier for this stream.

This refers to the external source system's account/tenant identifier (e.g., Snowflake account), not the Domo account object.

Returns:

Type Description
str | None

Source system account identifier or None if unavailable.

Example

stream.source_account "my_snowflake_account"

spreadsheet property

spreadsheet: str | None

Get spreadsheet identifier from stream configuration.

Works for Google connectors: - Google Sheets - Google Spreadsheets

Returns:

Type Description
str | None

Spreadsheet ID/filename or None if not available

Example

stream.spreadsheet "1A2B3C4D5E6F7G8H9I"

sql property

sql: str | None

Get SQL query from stream configuration (cross-platform).

Works across providers that use SQL queries: - Snowflake (all variants) - AWS Athena - Amazon Athena High Bandwidth - PostgreSQL

Returns:

Type Description
str | None

SQL query string or None if not available

Example

stream = await DomoStream.get_by_id(auth, stream_id) print(stream.sql) "SELECT * FROM my_table"

table property

table: str | None

Get table name from stream configuration (cross-platform).

Works for providers that operate on specific tables: - AWS Athena - Snowflake Writeback

Returns:

Type Description
str | None

Table name or None if not available

Example

stream.table "my_table"

typed_config property

typed_config

Get typed stream configuration built from raw API config.

Returns the appropriate typed config class based on data_provider_key and the raw configuration list from the API response.

warehouse property

warehouse: str | None

Get warehouse/compute resource from stream configuration.

Works for Snowflake variants that use compute warehouses.

Returns:

Type Description
str | None

Warehouse name or None if not available

Example

stream.warehouse "COMPUTE_WH"

extract_schedule_from_raw

extract_schedule_from_raw()

Extract schedule from stream configuration if available

Source code in src/crew_dcs/classes/DomoDataset/stream.py
530
531
532
533
534
535
536
def extract_schedule_from_raw(self):
    """Extract schedule from stream configuration if available"""

    if self.raw:
        self.Schedule = DomoSchedule.from_parent(parent=self, obj=self.raw)

    return self.Schedule

generate_update_body

generate_update_body(
    *,
    data_provider_key: str | None = None,
    transport_description: str | None = None,
    configuration: list[dict] | None = None,
    account_id: int | None = None,
    dataset_name: str | None = None,
    dataset_description: str | None = None,
    update_method: str | None = None
) -> dict

Generate the API body for PUT /api/data/v1/streams/{id}.

Composes a valid stream update payload from the current stream state, with optional overrides for any field. This is the single source of truth for how to build stream update bodies.

Parameters:

Name Type Description Default
data_provider_key str | None

Override the dataProvider.key (connector type). Defaults to current stream's data_provider_key.

None
transport_description str | None

Override the transport.description (e.g. 'com.domo.connector.domostats'). Defaults to current.

None
configuration list[dict] | None

Override the stream configuration list. Defaults to current stream's configuration from raw API response.

None
account_id int | None

Optional account ID for connectors that require auth.

None
dataset_name str | None

Override the dataSource.name. Defaults to parent dataset's name.

None
dataset_description str | None

Override the dataSource.description. Defaults to parent dataset's description.

None
update_method str | None

Override the updateMethod (e.g. 'REPLACE', 'APPEND'). Defaults to current stream's update_method.

None

Returns:

Name Type Description
dict dict

Valid body for PUT /api/data/v1/streams/{id}

Example

stream = await DomoStream.get_by_id(auth, stream_id="123") body = stream.generate_update_body( ... data_provider_key="domostats", ... transport_description="com.domo.connector.domostats", ... configuration=[{"name": "report", "value": "cards", "type": "string"}], ... ) await stream.update(body)

Source code in src/crew_dcs/classes/DomoDataset/stream.py
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
def generate_update_body(
    self,
    *,
    data_provider_key: str | None = None,
    transport_description: str | None = None,
    configuration: list[dict] | None = None,
    account_id: int | None = None,
    dataset_name: str | None = None,
    dataset_description: str | None = None,
    update_method: str | None = None,
) -> dict:
    """Generate the API body for PUT /api/data/v1/streams/{id}.

    Composes a valid stream update payload from the current stream state,
    with optional overrides for any field. This is the single source of truth
    for how to build stream update bodies.

    Args:
        data_provider_key: Override the dataProvider.key (connector type).
            Defaults to current stream's data_provider_key.
        transport_description: Override the transport.description
            (e.g. 'com.domo.connector.domostats'). Defaults to current.
        configuration: Override the stream configuration list.
            Defaults to current stream's configuration from raw API response.
        account_id: Optional account ID for connectors that require auth.
        dataset_name: Override the dataSource.name. Defaults to parent
            dataset's name.
        dataset_description: Override the dataSource.description.
            Defaults to parent dataset's description.
        update_method: Override the updateMethod (e.g. 'REPLACE', 'APPEND').
            Defaults to current stream's update_method.

    Returns:
        dict: Valid body for PUT /api/data/v1/streams/{id}

    Example:
        >>> stream = await DomoStream.get_by_id(auth, stream_id="123")
        >>> body = stream.generate_update_body(
        ...     data_provider_key="domostats",
        ...     transport_description="com.domo.connector.domostats",
        ...     configuration=[{"name": "report", "value": "cards", "type": "string"}],
        ... )
        >>> await stream.update(body)
    """
    # Resolve overrides with current state as fallback
    dp_key = data_provider_key or self.data_provider_key
    transport_desc = transport_description or self.transport_description
    method = update_method or self.update_method

    # Configuration: use override, or extract from raw API response
    # Internal config keys that cause 500 errors on create/update
    _skip_config_keys = {
        "retry.retryNumber",
        "updatemode.mode",
        "_description_",
    }

    if configuration is not None:
        config = [
            c
            for c in configuration
            if c.get("name") and c.get("name") not in _skip_config_keys
        ]
    else:
        raw_config = self.raw.get("configuration", []) if self.raw else []
        # Strip streamId and internal keys from config items
        config = [
            {
                "name": c.get("name"),
                "value": c.get("value"),
                "type": c.get("type", "string"),
                **({"category": c["category"]} if "category" in c else {}),
            }
            for c in raw_config
            if c.get("name") and c.get("name") not in _skip_config_keys
        ]

    # Dataset name/description from parent or override
    parent = self.parent
    ds_name = (
        dataset_name or (getattr(parent, "name", None) if parent else None) or ""
    )
    ds_desc = (
        dataset_description
        or (getattr(parent, "description", None) if parent else None)
        or ""
    )

    body: dict[str, Any] = {
        "transport": {
            "type": "CONNECTOR",
            "description": transport_desc,
            "version": str(self.transport_version or "0"),
        },
        "configuration": config,
        "updateMethod": method,
        "dataProvider": {"key": dp_key},
        "dataSource": {
            "name": ds_name,
            "description": ds_desc,
        },
    }

    if account_id is not None:
        body["account"] = {"id": account_id}

    return body

get_account async

get_account(
    session: AsyncClient | None = None,
    debug_api: bool = False,
    force_refresh: bool = False,
    is_suppress_no_account_config: bool = True,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> Any | None

Retrieve the Account associated with this stream.

Parameters:

Name Type Description Default
session AsyncClient | None

HTTP client session

None
debug_api bool

Enable API debugging

False
force_refresh bool

If True, refresh even if Account is already set

False
is_suppress_no_account_config bool

If True, suppress errors when account config is not found

True
context RouteContext | None

Optional RouteContext for API call configuration

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description
Any | None

DomoAccount instance or None if no account_id

Example

stream = await DomoStream.get_by_id(auth=auth, stream_id="123") account = await stream.get_account() print(f"Account: {account.name}")

Source code in src/crew_dcs/classes/DomoDataset/stream.py
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
async def get_account(
    self,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    force_refresh: bool = False,
    is_suppress_no_account_config: bool = True,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> Any | None:  # DomoAccount
    """Retrieve the Account associated with this stream.

    Args:
        session: HTTP client session
        debug_api: Enable API debugging
        force_refresh: If True, refresh even if Account is already set
        is_suppress_no_account_config: If True, suppress errors when account config is not found
        context: Optional RouteContext for API call configuration
        **context_kwargs: Additional context parameters

    Returns:
        DomoAccount instance or None if no account_id

    Example:
        >>> stream = await DomoStream.get_by_id(auth=auth, stream_id="123")
        >>> account = await stream.get_account()
        >>> print(f"Account: {account.name}")
    """
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        **context_kwargs,
    )

    if not self.account_id:
        return None

    if self.Account is not None and not force_refresh:
        return self.Account

    from ..DomoAccount import DomoAccount

    try:
        self.Account = await DomoAccount.get_by_id(
            auth=self.auth,
            account_id=self.account_id,
            context=context,
            is_use_default_account_class=False,
            is_suppress_no_config=is_suppress_no_account_config,
        )
    except dmde.DomoError as e:
        if is_suppress_no_account_config:
            await logger.warning(
                f"Warning: Could not retrieve account {self.account_id}: {e}"
            )
            self.Account = None
        else:
            raise

    return self.Account

get_by_id async classmethod

get_by_id(
    auth: DomoAuth,
    stream_id: str,
    return_raw: bool = False,
    debug_num_stacks_to_drop=2,
    debug_api: bool = False,
    session: AsyncClient | None = None,
    is_get_account: bool = True,
    is_suppress_no_account_config: bool = True,
    *,
    context: RouteContext | None = None,
    **context_kwargs
)

Get a stream by its ID.

Parameters:

Name Type Description Default
auth DomoAuth

Authentication object

required
stream_id str

Unique stream identifier

required
return_raw bool

Return raw response without processing

False
debug_num_stacks_to_drop

Stack frames to drop for debugging

2
debug_api bool

Enable API debugging

False
session AsyncClient | None

HTTP client session

None
is_get_account bool

If True and account_id is present, retrieve full Account object

True
is_suppress_no_account_config bool

If True, suppress errors when account config is not found

True
context RouteContext | None

Optional RouteContext for API call configuration

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description

DomoStream instance or ResponseGetData if return_raw=True

Raises:

Type Description
Stream_GET_Error

If stream retrieval fails

Source code in src/crew_dcs/classes/DomoDataset/stream.py
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
@classmethod
async def get_by_id(
    cls,
    auth: DomoAuth,
    stream_id: str,
    return_raw: bool = False,
    debug_num_stacks_to_drop=2,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    is_get_account: bool = True,
    is_suppress_no_account_config: bool = True,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Get a stream by its ID.

    Args:
        auth: Authentication object
        stream_id: Unique stream identifier
        return_raw: Return raw response without processing
        debug_num_stacks_to_drop: Stack frames to drop for debugging
        debug_api: Enable API debugging
        session: HTTP client session
        is_get_account: If True and account_id is present, retrieve full Account object
        is_suppress_no_account_config: If True, suppress errors when account config is not found
        context: Optional RouteContext for API call configuration
        **context_kwargs: Additional context parameters

    Returns:
        DomoStream instance or ResponseGetData if return_raw=True

    Raises:
        Stream_GET_Error: If stream retrieval fails
    """

    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        debug_num_stacks_to_drop=debug_num_stacks_to_drop,
        **context_kwargs,
    )

    res = await stream_routes.get_stream_by_id(
        auth=auth,
        stream_id=stream_id,
        context=context,
    )

    if return_raw:
        return res

    stream = cls.from_dict(auth=auth, obj=res.response)

    # Retrieve Account if account_id is present
    if is_get_account and stream.account_id:
        await stream.get_account(
            context=context,
            force_refresh=True,
            is_suppress_no_account_config=is_suppress_no_account_config,
        )
    return stream

swap_connector async

swap_connector(
    data_provider_key: str,
    transport_description: str,
    configuration: list[dict] | None = None,
    *,
    account_id: int | None = None,
    dataset_name: str | None = None,
    dataset_description: str | None = None,
    update_method: str | None = None,
    is_refresh: bool = True,
    context: RouteContext | None = None,
    **context_kwargs
) -> DomoStream

Swap this stream's connector type.

Changes the dataProvider.key and associated configuration, effectively changing how data gets INTO the dataset. The dataset keeps the same dataset_id and stream_id — only the connector changes.

This is the core mechanism for: - Changing a dataset's connector type after creation - Restoring a deleted connector's dataset with a different connector - Migrating data sources between connector types

Parameters:

Name Type Description Default
data_provider_key str

The new dataProvider.key (connector type identifier). e.g. 'domostats', 'snowflake', 'domo-governance-...'

required
transport_description str

The new transport description (connector class path). e.g. 'com.domo.connector.domostats'

required
configuration list[dict] | None

The new stream configuration list. e.g. [{"name": "report", "value": "cards", "type": "string"}]

None
account_id int | None

Optional account ID for connectors that require auth. Some connectors (e.g. governance, Snowflake) need an account with credentials. Without this, execution will fail with AccountAssociationException.

None
dataset_name str | None

Override the dataset name. Defaults to current name.

None
dataset_description str | None

Override the dataset description.

None
update_method str | None

Override the updateMethod. Defaults to current.

None
is_refresh bool

If True, refresh the stream object from the API response after the swap. Default True.

True
context RouteContext | None

Optional RouteContext for API call configuration.

None
**context_kwargs

Additional context parameters.

{}

Returns:

Name Type Description
DomoStream DomoStream

self, updated with the new connector state.

Raises:

Type Description
Stream_CRUD_Error

If the swap fails.

Example

stream = await DomoStream.get_by_id(auth, stream_id="123")

Swap from governance to domostats

await stream.swap_connector( ... data_provider_key="domostats", ... transport_description="com.domo.connector.domostats", ... configuration=[ ... {"name": "report", "value": "cards", "type": "string"}, ... ], ... )

Run the stream with the new connector

await stream_routes.execute_stream(auth=auth, stream_id=stream.id)

Note

The dataset's dataProviderType metadata may still show the original connector type after the swap. The stream's dataProvider.key is what matters for execution.

Source code in src/crew_dcs/classes/DomoDataset/stream.py
 909
 910
 911
 912
 913
 914
 915
 916
 917
 918
 919
 920
 921
 922
 923
 924
 925
 926
 927
 928
 929
 930
 931
 932
 933
 934
 935
 936
 937
 938
 939
 940
 941
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
async def swap_connector(
    self,
    data_provider_key: str,
    transport_description: str,
    configuration: list[dict] | None = None,
    *,
    account_id: int | None = None,
    dataset_name: str | None = None,
    dataset_description: str | None = None,
    update_method: str | None = None,
    is_refresh: bool = True,
    context: RouteContext | None = None,
    **context_kwargs,
) -> DomoStream:
    """Swap this stream's connector type.

    Changes the dataProvider.key and associated configuration, effectively
    changing how data gets INTO the dataset. The dataset keeps the same
    dataset_id and stream_id — only the connector changes.

    This is the core mechanism for:
    - Changing a dataset's connector type after creation
    - Restoring a deleted connector's dataset with a different connector
    - Migrating data sources between connector types

    Args:
        data_provider_key: The new dataProvider.key (connector type identifier).
            e.g. 'domostats', 'snowflake', 'domo-governance-...'
        transport_description: The new transport description
            (connector class path). e.g. 'com.domo.connector.domostats'
        configuration: The new stream configuration list.
            e.g. [{"name": "report", "value": "cards", "type": "string"}]
        account_id: Optional account ID for connectors that require auth.
            Some connectors (e.g. governance, Snowflake) need an account
            with credentials. Without this, execution will fail with
            AccountAssociationException.
        dataset_name: Override the dataset name. Defaults to current name.
        dataset_description: Override the dataset description.
        update_method: Override the updateMethod. Defaults to current.
        is_refresh: If True, refresh the stream object from the API
            response after the swap. Default True.
        context: Optional RouteContext for API call configuration.
        **context_kwargs: Additional context parameters.

    Returns:
        DomoStream: self, updated with the new connector state.

    Raises:
        Stream_CRUD_Error: If the swap fails.

    Example:
        >>> stream = await DomoStream.get_by_id(auth, stream_id="123")
        >>> # Swap from governance to domostats
        >>> await stream.swap_connector(
        ...     data_provider_key="domostats",
        ...     transport_description="com.domo.connector.domostats",
        ...     configuration=[
        ...         {"name": "report", "value": "cards", "type": "string"},
        ...     ],
        ... )
        >>> # Run the stream with the new connector
        >>> await stream_routes.execute_stream(auth=auth, stream_id=stream.id)

    Note:
        The dataset's `dataProviderType` metadata may still show the
        original connector type after the swap. The stream's
        `dataProvider.key` is what matters for execution.
    """
    body = self.generate_update_body(
        data_provider_key=data_provider_key,
        transport_description=transport_description,
        configuration=configuration,
        account_id=account_id,
        dataset_name=dataset_name,
        dataset_description=dataset_description,
        update_method=update_method,
    )

    context = RouteContext.build_context(context=context, **context_kwargs)

    res = await self.update(body, context=context)

    if is_refresh and res.is_success:
        # Refresh self from the response
        updated = DomoStream.from_dict(
            auth=self.auth,
            obj=res.response,
            parent=self.parent,
        )
        # Copy updated fields back to self
        self.data_provider_key = updated.data_provider_key
        self.data_provider_name = updated.data_provider_name
        self.transport_description = updated.transport_description
        self.transport_version = updated.transport_version
        self.update_method = updated.update_method
        self.account_id = updated.account_id
        self.raw = updated.raw
        self._typed_config = None  # Reset typed config for new provider

    return self

to_dict

to_dict(
    override_fn=None,
    return_snake_case: bool = False,
    *,
    include_configs: bool = True
) -> dict

Serialize stream and related configuration for export/debugging.

Source code in src/crew_dcs/classes/DomoDataset/stream.py
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
def to_dict(
    self,
    override_fn=None,
    return_snake_case: bool = False,
    *,
    include_configs: bool = True,
) -> dict:
    """Serialize stream and related configuration for export/debugging."""
    if override_fn:
        return override_fn(self)

    result = super().to_dict(return_snake_case=return_snake_case)

    if include_configs:
        if self.typed_config:
            key = "typed_config" if return_snake_case else "TypedConfig"
            result[key] = self.typed_config.to_dict(
                return_snake_case=return_snake_case,
                for_api=False,
            )

        if self.Account and getattr(self.Account, "Config", None):
            key = "account_config" if return_snake_case else "AccountConfig"
            result[key] = self.Account.Config.to_dict(
                return_snake_case=return_snake_case
            )

    return result

Stream_CRUD_Error

Stream_CRUD_Error(
    operation: str,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs
)

Bases: RouteError

Raised when stream create, update, delete, or execute operations fail.

Source code in src/crew_dcs/routes/dataset/stream.py
61
62
63
64
65
66
67
68
69
70
71
72
73
74
def __init__(
    self,
    operation: str,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs,
):
    super().__init__(
        message=message or f"Stream {operation} operation failed",
        entity_id=stream_id,
        res=res,
        **kwargs,
    )

Stream_GET_Error

Stream_GET_Error(
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs
)

Bases: RouteError

Raised when stream retrieval operations fail.

Source code in src/crew_dcs/routes/dataset/stream.py
43
44
45
46
47
48
49
50
51
52
53
54
55
def __init__(
    self,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs,
):
    super().__init__(
        message=message or "Stream retrieval failed",
        entity_id=stream_id,
        res=res,
        **kwargs,
    )