Manage Workflow Assets, Tables, and Dataset Lineage

For the complete documentation index, see llms.txt. For a full content snapshot, see llms-full.txt. Append .md to any kestra.io/docs/* URL for plain Markdown.

Available on:Enterprise EditionCloud

Track and manage the resources your workflows create and use.

An asset is any named resource a workflow reads from or writes to — a database table, a file, a virtual machine. Declaring assets on tasks builds a lineage graph: which workflows touch which resources, in what order, and how they depend on each other. Kestra can ship that graph to external lineage platforms such as DataHub, Marquez, or Atlan via OpenLineage, so orchestration lineage appears alongside warehouse and pipeline lineage in one place.

Assets are captured automatically when tasks declare assets.inputs or assets.outputs. You can also add them manually from the Assets tab. Assets enable:

  • Shipping metadata to lineage providers (e.g., OpenLineage).
  • Populating dropdowns or Pebble inputs with live assets (e.g., available VMs).
  • Monitoring assets and their state.

Asset definition

Define assets directly on any task using the assets property. Each task can declare inputs assets (resources it reads) and outputs assets (resources it creates or modifies).

Every asset includes these fields:

FieldDescription
idunique within a tenant
namespaceassociates the asset with a namespace for filtering and RBAC management. Set to null to create a Global asset visible to all namespaces in a tenant.
typeuse predefined Kestra types like io.kestra.plugin.ee.assets.Table or any custom string value
displayNameoptional human-readable name
descriptionmarkdown-supported documentation
metadatamap of key-value for adding custom metadata to the given asset

Asset identifier

An asset is uniquely identified by its id and the tenant (tenantId) where you create it - the id must be unique per tenant. Neither the namespace nor the type is part of that identity: two assets with the same id and different namespaces or types cannot exist in the same tenant. Creating an asset with an id that is already taken is rejected.

You can attach a namespace to an asset to improve filtering and to restrict visibility so only users or groups with the appropriate RBAC can access the asset. The namespace field is editable; you can change it from the UI or by updating the namespace field in a flow’s assets.outputs declaration.

An asset that doesn’t have any namespace set (i.e., the namespace field is empty) can be considered a global asset. Global assets are not guarded by namespace-level ACL: any user or group with a namespace-scoped grant can see and interact with them. Use them for shared resources that every namespace in the tenant should be able to reference, such as a central data catalog entry or a shared infrastructure asset.

To create a Global asset from a flow, omit the namespace field on the assets.outputs entry entirely. A flow has no way to pass an explicit null through YAML, so omitting the field on a new asset is the only way to produce a null-namespace (Global) asset from a flow definition. When updating an existing asset, omitting namespace preserves the asset’s current namespace rather than clearing it.

Cross-namespace asset writes

When a flow writes to an asset in a different namespace, Kestra checks that the flow’s own namespace is allowed to write to the target namespace. If the check fails, the write is rejected. This applies to both assets.outputs writes and assets.inputs retype operations.

If you have flows that write to assets across namespace boundaries, ensure the flow’s namespace is listed in the target namespace’s allowedNamespaces configuration. Flows that previously bypassed this check may start failing after upgrading.

Asset type

Asset types fall into two categories:

  • Kestra-defined asset types: These predefined types use the io.kestra.core.models.assets model and provide structured metadata fields specific to each asset type. Plugins that support auto-generation populate these fields automatically during task execution — for example, a JDBC plugin creates a Table asset with system, database, and schema filled in from the connection details.

Kestra provides these built-in asset types:

  • io.kestra.plugin.ee.assets.Dataset

    • Represents a dataset asset managed by Kestra.
    • Metadata: system, location, format
  • io.kestra.plugin.ee.assets.File

    • Represents a file asset, such as documents, logs, or other file-based outputs.
    • Metadata: system, path
  • io.kestra.plugin.ee.assets.Table

    • Represents a database table asset with schema and data location metadata.
    • Metadata: system, database, schema
  • io.kestra.plugin.ee.assets.VM

    • Represents a virtual machine asset, including attributes like IP address and provider.
    • Metadata: provider, region, state
  • io.kestra.core.models.assets.External

    • Represents an external asset that exists outside of Kestra’s managed resources.
    • External is the placeholder type assigned when you reference an asset in assets.inputs with no declared type and the asset does not yet exist. If you declare a type on the assets.inputs entry, Kestra uses that declared type when creating the asset and will retype the asset to match if it was previously created as External. You do not need to set type explicitly unless you want the asset to carry a specific type.
    • This is useful for tracking dependencies on resources managed outside your workflows, such as external database tables, third-party APIs, or manually provisioned infrastructure.
  • Free-form asset types: You can define asset types using any custom string value to represent asset categories that fit your organization’s needs. This lets you create and manage your own asset taxonomies, giving you flexibility to describe resources that are not covered by Kestra’s standard models. These assets require manual definition and will not be auto-generated by plugins.

Quick start: minimal asset flow

A small example that registers one output asset and logs its ID:

id: hello_assets
namespace: company.team
tasks:
- id: write_file
type: io.kestra.plugin.core.log.Log
message: "Created report.csv"
assets:
outputs:
- id: report.csv
type: io.kestra.plugin.ee.assets.File
metadata:
path: s3://company/reports/report.csv
- id: confirm
type: io.kestra.plugin.core.log.Log
message: "Asset recorded: {{ assets() | jq('.[] | {id: .id, type: .type, metadata: .metadata}') }}"

Auto-generated assets

Some plugins support automatic asset generation when assets.enableAuto: true is set on a task. This removes the need to manually declare assets.inputs and assets.outputs — the plugin inspects its execution context and emits assets automatically:

  • JDBC Query: detects CREATE TABLE statements and emits a single io.kestra.plugin.ee.assets.Table output; JDBC URL populates system and database.

  • Ansible CLI: parses inventory hosts as inputs of type io.kestra.core.models.assets.External, marking the infrastructure targets the playbook runs against.

  • dbt CLI and dbt Cloud CheckStatus: parse manifest.json and run_results.json to emit models, seeds, and snapshots as io.kestra.plugin.ee.assets.Table outputs with database, schema, name, and lineage edges based on depends_on. Test outcomes ride the model’s asset as metadata: dbtTestStatus (pass, fail, warn), dbtTestsTotal, and dbtTestsFailed. A relationships test rolls onto every model it targets. Skipped tests carry no keys; a model with no executed tests carries no keys either.

  • Helm: Upgrade, Rollback, Status, and Uninstall emit a Custom output typed io.kestra.plugin.ee.assets.HelmRelease for the release plus one Custom output typed io.kestra.plugin.ee.assets.KubernetesResource per managed resource (Deployment, Service, Ingress, etc.) — these aren’t typed asset classes yet, but the type strings are chosen to match what they’d become if core-ee adds them later, so nothing in the catalog needs to change on that swap; the chart reference and any valuesFrom files are declared as inputs.

  • Qlik Cloud apps.Reload: emits a Custom asset typed io.kestra.plugin.qlikcloud.assets.App for the reloaded app, carrying the app’s freshness. The asset id is the Qlik app id. Declare assets.inputs manually to connect upstream dbt or Fivetran assets to this node.

  • Hex projects.Run: emits a Custom asset typed io.kestra.plugin.ee.assets.Dataset (with system: hex) for the Hex project that ran, so Hex appears as the terminal consumer in a Fivetran → dbt → Hex lineage chain. The asset id is the projectId. Hex’s API reports no upstream tables, so declare assets.inputs manually using the same database.schema.table ids that plugin-dbt and plugin-fivetran emit.

JDBC Query auto-generated assets
id: jdbc_create_trips
namespace: company.team
tasks:
- id: create_trips_table
type: io.kestra.plugin.jdbc.sqlite.Query
url: jdbc:sqlite:myfile.db
outputDbFile: true
sql: |
CREATE TABLE IF NOT EXISTS trips (
VendorID INTEGER,
passenger_count INTEGER,
trip_distance REAL
);
assets:
enableAuto: true
Ansible CLI auto-generated assets
id: ansible_playbook
namespace: company.team
tasks:
- id: ansible_task
type: io.kestra.plugin.ansible.cli.AnsibleCLI
inputFiles:
inventory.ini: |
localhost ansible_connection=local
myplaybook.yml: |
---
- hosts: localhost
tasks:
- name: Print Hello World
debug:
msg: "Hello, World!"
assets:
enableAuto: true
commands:
- ansible-playbook -i inventory.ini myplaybook.yml
dbt CLI auto-generated assets
id: dbt_build_duckdb
namespace: company.team
tasks:
- id: dbt
type: io.kestra.plugin.core.flow.WorkingDirectory
tasks:
- id: clone_repository
type: io.kestra.plugin.git.Clone
url: https://github.com/kestra-io/dbt-example
branch: main
- id: dbt_build
type: io.kestra.plugin.dbt.cli.DbtCLI
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
containerImage: ghcr.io/kestra-io/dbt-duckdb:latest
commands:
- dbt deps
- dbt build
- dbt run
profiles: |
my_dbt_project:
outputs:
dev:
type: duckdb
path: ":memory:"
fixed_retries: 1
threads: 16
timeout_seconds: 300
target: dev
assets:
enableAuto: true
Helm auto-generated assets
id: helm_upgrade_release
namespace: company.team
tasks:
- id: upgrade
type: io.kestra.plugin.helm.Upgrade
releaseName: nginx
namespace: web
chart:
repository: https://kubernetes.github.io/ingress-nginx
name: ingress-nginx
version: 4.11.3
assets:
enableAuto: true
Qlik Cloud app reload auto-generated assets
id: qlik_reload
namespace: company.team
tasks:
- id: reload
type: io.kestra.plugin.qlikcloud.apps.Reload
tenantUrl: https://mytenant.eu.qlikcloud.com
apiKey: "{{ secret('QLIK_API_KEY') }}"
appId: 60f2e3b1a1b2c3d4e5f6a7b8
assets:
enableAuto: true
inputs:
- id: analytics.marts.fct_orders
type: io.kestra.plugin.ee.assets.Table
Hex project run auto-generated assets
id: hex_run_project
namespace: company.team
tasks:
- id: run
type: io.kestra.plugin.hex.projects.Run
apiKey: "{{ secret('HEX_API_KEY') }}"
projectId: 60f2e3b1-a1b2-c3d4-e5f6-a7b8c9d0e1f2
assets:
enableAuto: true
inputs:
- id: analytics.marts.fct_orders
type: io.kestra.plugin.ee.assets.Table

Operational automation

Assets also support lifecycle management, event-driven triggers, and freshness monitoring directly from flows:

  • Lifecycle tasks to create, update, list, and delete assets (Set, List, Delete).
  • Event-based triggers with EventTrigger that react to asset lifecycle events (CREATED, UPDATED, DELETED, USED).
  • Freshness monitoring with FreshnessTrigger to detect stale assets and launch flows automatically.
  • Scope triggers by asset ID, namespace, type, and metadata filters.
  • Trigger context variables (event, eventTime, lastUpdated, staleDuration, checkTime) available for routing, alerting, and recovery logic.

Freshness states for each asset are visible in the dependencies graph.

Trigger use mapping

TriggerPrimary use
EventTriggerReact instantly to asset lifecycle events (CREATED, UPDATED, DELETED, USED).
FreshnessTriggerPoll assets on an interval to detect staleness and launch remediation.
Advanced: event-driven automation
id: asset_event_driven_pipeline
namespace: company.data
tasks:
- id: transform_to_mart
type: io.kestra.plugin.core.flow.Subflow
namespace: company.data
flowId: create_mart_tables
inputs:
source_asset_id: "{{ trigger.asset.id }}"
source_event: "{{ trigger.asset.event }}"
event_time: "{{ trigger.asset.eventTime }}"
triggers:
- id: staging_table_event
type: io.kestra.plugin.ee.assets.EventTrigger
namespace: company.data
assetType: io.kestra.plugin.ee.assets.Table
events:
- CREATED
- UPDATED
metadataQuery:
- field: model_layer
type: EQUAL_TO
value: staging
Advanced: audit deletions
id: audit_asset_deletions
namespace: company.security
tasks:
- id: log_deletion
type: io.kestra.plugin.jdbc.postgresql.Query
sql: |
INSERT INTO audit_log (asset_id, asset_type, namespace, event, event_time)
VALUES (
'{{ trigger.asset.id }}',
'{{ trigger.asset.type }}',
'{{ trigger.asset.namespace }}',
'{{ trigger.asset.event }}',
'{{ trigger.asset.eventTime }}'
)
triggers:
- id: asset_deletion_event
type: io.kestra.plugin.ee.assets.EventTrigger
events:
- DELETED
Advanced: alert on dbt test failures
id: dbt_test_failure_alert
namespace: company.data
tasks:
- id: notify
type: io.kestra.plugin.core.log.Log
message: >
dbt tests failed on asset {{ trigger.asset.id }}.
Status: {{ trigger.asset.metadata.dbtTestStatus }},
failed: {{ trigger.asset.metadata.dbtTestsFailed }} of {{ trigger.asset.metadata.dbtTestsTotal }}.
triggers:
- id: dbt_tests_failed
type: io.kestra.plugin.ee.assets.EventTrigger
events:
- UPDATED
metadataQuery:
- field: dbtTestStatus
type: EQUAL_TO
value: fail
Advanced: freshness monitoring
id: stale_assets_monitor
namespace: company.monitoring
tasks:
- id: log_stale
type: io.kestra.plugin.core.log.Log
message: >
Found {{ trigger.assets | length }} stale assets.
First asset: {{ trigger.asset.id ?? 'n/a' }}.
Stale for: {{ trigger.asset.staleDuration ?? 'n/a' }}.
triggers:
- id: stale_assets
type: io.kestra.plugin.ee.assets.FreshnessTrigger
maxStaleness: PT24H
interval: PT1H
Advanced: scoped freshness checks
id: prod_assets_freshness
namespace: company.monitoring
tasks:
- id: trigger_remediation
type: io.kestra.plugin.core.flow.Subflow
namespace: company.data
flowId: refresh_marts
inputs:
asset_id: "{{ trigger.asset.id }}"
last_updated: "{{ trigger.asset.lastUpdated }}"
stale_duration: "{{ trigger.asset.staleDuration }}"
triggers:
- id: stale_prod_marts
type: io.kestra.plugin.ee.assets.FreshnessTrigger
namespace: company.data
assetType: TABLE
maxStaleness: PT6H
interval: PT30M
metadataQuery:
- field: environment
type: EQUAL_TO
value: prod
- field: model_layer
type: EQUAL_TO
value: mart
Advanced: lifecycle tasks
id: asset_lifecycle_ops
namespace: company.data
tasks:
- id: upsert_asset
type: io.kestra.plugin.ee.assets.Set
namespace: assets.data
assetId: customers_by_country
assetType: TABLE
displayName: Customers by Country
assetDescription: Customer distribution by country
metadata:
owner: data-team
environment: prod
- id: list_assets
type: io.kestra.plugin.ee.assets.List
namespace: assets.data
types:
- TABLE
metadataQuery:
- field: owner
type: EQUAL_TO
value: data-team
fetchType: FETCH
- id: delete_asset
type: io.kestra.plugin.ee.assets.Delete
assetId: customers_by_country

Explore the dependencies graph

The dependencies graph shows how assets relate to each other across your workflows. Open it from the Graph tab.

Tree and DAG layouts

Two layout modes are available via the toggle in the graph toolbar:

  • DAG (directed acyclic graph): a deterministic ranked layout. Assets are ordered left-to-right by longest path, so the graph looks identical on every reload. Use this for a stable, presentation-ready view.
  • Tree: a force-directed layout that distributes nodes more evenly across the canvas. Use this for dense graphs where the DAG layout produces overlaps.

Both modes show the same nodes and edges.

Freshness

Each node displays the freshness state of that asset based on the most recent execution that produced it:

StateMeaning
freshThe asset was produced successfully within the expected cadence.
staleThe asset has not been updated within the expected cadence.
failedThe most recent producing execution failed.
unknownNo producing execution has been recorded for this asset.

A summary bar above the graph shows the count of each state; the legend shows only states present in the graph.

The expected cadence is derived from the Schedule trigger of the producing flow. To launch remediation flows when an asset becomes stale, use FreshnessTrigger.

Group by

Use the Group by selector to cluster nodes into labeled buckets:

  • dataset: groups assets by the dataset field in the asset’s schema metadata.
  • producer: groups assets by the plugin artifact: the fourth segment of the producing task type’s FQCN (e.g., jdbc, dbt, aws).

Selected groups appear as chips in a row below the toolbar. Hovering a chip fades unrelated nodes; clicking pins the group so it stays highlighted. Click the chip again or click an empty area of the canvas to release it. Flow nodes appear in their own bucket and are not merged into asset groups.

Node details panel

Click a node to open the details panel on the right side of the graph. The panel shows:

  • The asset’s full identifier and type.
  • Its current status and metadata.
  • Recent runs, each linked to its execution page.

Double-click a node to navigate to that asset’s detail page. Clicking an empty area of the canvas closes the panel and releases any pinned group.

Locking assets

A lock prevents concurrent writes to a shared asset while a flow operates on it. Locks are TTL-bounded: they expire automatically when their duration elapses and can also be released explicitly. Reads are always open — only writes (edit, delete) are blocked while a lock is held.

Two owner types exist:

Owner typeAcquired byBehavior while held
EXECUTIONAcquire taskBlocks other executions’ writes. The lock-holding execution can still write to the asset, and each write extends the lease.
USERUI or REST APIBlocks all execution writes. Use for manual maintenance windows.

Locking from a flow

The Acquire and Release tasks wrap the work that needs exclusive write access. Both require the LOCK permission on the ASSET resource (UNLOCK for Release).

If another execution already holds the lock when Acquire runs, the task fails with a 423 error. Add a Retry to the Acquire task to wait for the lock to become available.

id: update_customer_asset
namespace: company.team
tasks:
- id: acquire
type: io.kestra.plugin.kestra.ee.locks.Acquire
assetId: customers_by_country
ttl: PT1H
- id: write
type: io.kestra.plugin.core.log.Log
message: Writing to the locked asset
- id: release
type: io.kestra.plugin.kestra.ee.locks.Release
assetId: customers_by_country

Acquire properties:

PropertyRequiredDescription
assetIdYesID of the asset to lock.
ttlNoHow long to hold the lock before it expires automatically. ISO-8601 duration (e.g. PT1H). Defaults to 5 minutes when unset.

Acquire outputs:

OutputDescription
lockedUntilWhen the lock expires.
ownerTypeAlways EXECUTION for a task-acquired lock.
executionIdID of the execution holding the lock.

Release is owner-checked: it removes the lock only if the current execution holds it. If the lock has already expired or belongs to a different owner, Release is a no-op — safe to call unconditionally.

Locking from the UI

From any asset’s detail page, users with the LOCK permission can lock the asset manually. Choose from preset durations (5 minutes to 24 hours) or enter a custom ISO-8601 duration. The page shows who holds the lock and when it expires. Users with the UNLOCK permission can release any lock regardless of owner.

When an asset is locked, the detail page shows a banner: You might be seeing outdated metadata as this asset is currently locked for writing.

The asset list supports filtering by lock status.

Data pipeline use cases

Advanced: data pipeline examples

Assets are essential for tracking data lineage in analytics and data engineering workflows. The following examples demonstrate how to use assets for simple table creation and complex multi-layer data pipelines.

Example 1: Simple table creation

id: pipeline_with_assets
namespace: company.team
tasks:
- id: create_trips_table
type: io.kestra.plugin.jdbc.sqlite.Queries
url: jdbc:sqlite:myfile.db
outputDbFile: true
sql: |
CREATE TABLE IF NOT EXISTS trips (
VendorID INTEGER,
passenger_count INTEGER,
trip_distance REAL
);
INSERT INTO trips (VendorID, passenger_count, trip_distance) VALUES
(1, 1, 1.5),
(1, 2, 2.3),
(2, 1, 0.8),
(2, 3, 3.1);
assets:
outputs:
- id: trips
namespace: "{{ flow.namespace }}"
type: io.kestra.plugin.ee.assets.Table
metadata:
database: sqlite
table: trips

Example 2: Multi-layer data pipeline

The staging layer reads from an external source; the mart layer creates aggregated analytics tables.

id: data_pipeline_assets
namespace: kestra.company.data
tasks:
- id: create_staging_layer_asset
type: io.kestra.plugin.jdbc.duckdb.Query
url: "jdbc:duckdb:md:my_db?motherduck_token={{ secret('MOTHERDUCK_TOKEN') }}"
fetchType: STORE
sql: |
CREATE TABLE IF NOT EXISTS trips AS
select VendorID, passenger_count, trip_distance from sample_data.nyc.taxi limit 10;
assets:
inputs:
- id: sample_data.nyc.taxi
outputs:
- id: trips
namespace: "{{flow.namespace}}"
type: io.kestra.plugin.ee.assets.Table
metadata:
model_layer: staging
- id: for_each
type: io.kestra.plugin.core.flow.Loop
values:
- passenger_count
- trip_distance
tasks:
- id: create_mart_layer_asset
type: io.kestra.plugin.jdbc.duckdb.Query
url: "jdbc:duckdb:md:my_db?motherduck_token={{ secret('MOTHERDUCK_TOKEN') }}"
fetchType: STORE
sql: SELECT AVG({{item.value}}) AS avg_{{item.value}} FROM trips;
assets:
inputs:
- id: trips
outputs:
- id: avg_{{item.value}}
type: io.kestra.plugin.ee.assets.Table
namespace: "{{flow.namespace}}"
metadata:
model_layer: mart

Infrastructure use case: team bucket provisioning

Advanced: infrastructure provisioning

The following flow creates S3 buckets for selected teams and registers them as assets:

id: infra_assets
namespace: kestra.company.infra
inputs:
- id: teams
type: MULTISELECT
values:
- Business
- Data
- Finance
- Product
tasks:
- id: for_each
type: io.kestra.plugin.core.flow.Loop
values: "{{ inputs.teams }}"
tasks:
- id: create_bucket
type: io.kestra.plugin.aws.cli.AwsCLI
accessKeyId: "{{ secret('AWS_ACCESS_KEY') }}"
secretKeyId: "{{ secret('AWS_SECRET_ACCESS_KEY') }}"
region: "{{ secret('AWS_REGION') }}"
allowFailure: true
commands:
- aws s3 mb s3://kestra-{{ item.value | slugify }}-bucket
assets:
outputs:
- id: kestra-{{ item.value | slugify }}-bucket
type: AWS_BUCKET
metadata:
provider: s3
address: s3://kestra-{{ item.value | slugify }}-bucket

The flow dynamically creates buckets (e.g., kestra-data-bucket, kestra-finance-bucket) and registers each as an AWS_BUCKET asset with relevant metadata. Teams reference these assets in downstream workflows:

id: upload_file
namespace: kestra.company.data
tasks:
- id: download
type: io.kestra.plugin.core.http.Download
uri: https://huggingface.co/datasets/kestra/datasets/raw/main/jaffle-csv/raw_customers.csv
- id: aws_upload
type: io.kestra.plugin.aws.s3.Upload
accessKeyId: "{{ secret('AWS_ACCESS_KEY') }}"
secretKeyId: "{{ secret('AWS_SECRET_ACCESS_KEY') }}"
region: "{{ secret('AWS_REGION') }}"
bucket: kestra-data-bucket
from: '{{ outputs.download.uri }}'
key: raw_customer.csv
assets:
inputs:
- id: kestra-data-bucket
outputs:
- id: raw_customer
type: io.kestra.plugin.ee.assets.File
metadata:
owner: data

Populate dropdowns and app inputs

The assets() Pebble function queries assets at runtime, for example to populate dropdown inputs or select resources based on type, namespace, or metadata.

Function signature

assets(id: string, type: string, namespace: string, metadata: map)

Parameters

ParameterTypeRequiredDescription
idstringNoFilter by asset ID. Because IDs are unique per tenant, this returns at most one result.
typestringNoFilter by asset type (e.g., "io.kestra.plugin.ee.assets.Table"). If omitted, returns all types.
namespacestringNoFilter by namespace. Defaults to the flow’s namespace.
metadatamapNoFilter by metadata key-value pairs (e.g., {"model_layer": "mart"}).

Return value

Returns an array of asset objects. Each object contains:

FieldTypeDescription
idstringAsset identifier
namespacestringNamespace the asset belongs to
typestringAsset type
metadatamapCustom metadata key-value pairs
tenantIdstringTenant ID where the asset was created
createdstringISO 8601 timestamp of creation
updatedstringISO 8601 timestamp of last update
deletedbooleanWhether the asset has been deleted

Examples

Fetch a specific asset by ID:

id: check_asset
namespace: company.team
tasks:
- id: log
type: io.kestra.plugin.core.log.Log
message: "{{ assets(id='report.csv') | jq('.[0].metadata.path') }}"

Populate a multiselect dropdown with table assets:

id: select_assets
namespace: company.team
inputs:
- id: assets
type: MULTISELECT
expression: '{{ assets(type="io.kestra.plugin.ee.assets.Table") | jq(".[].id") }}'
tasks:
- id: for_each
type: io.kestra.plugin.core.flow.Loop
values: "{{inputs.assets}}"
tasks:
- id: log
type: io.kestra.plugin.core.log.Log
message: "{{item.value}}"

Filter assets by namespace:

inputs:
- id: staging_tables
type: MULTISELECT
expression: '{{ assets(type="io.kestra.plugin.ee.assets.Table", namespace="company.team") | jq(".[].id") }}'

Filter assets by metadata:

inputs:
- id: mart_tables
type: MULTISELECT
expression: '{{ assets(metadata={"model_layer": "mart"}) | jq(".[].id") }}'

Get all assets and extract metadata:

id: list_assets_metadata
namespace: company.team
tasks:
- id: list_all_assets
type: io.kestra.plugin.core.log.Log
message: "{{ assets() | jq('.[] | {id: .id, type: .type, metadata: .metadata}') }}"

Export assets with AssetShipper

The AssetShipper task exports asset metadata to external systems for lineage tracking, monitoring, or integration with data catalogs. Supported destinations include files and OpenLineage-compatible providers.

Export assets to file

Export asset metadata to a file in either ION or JSON format. This is useful for archiving, auditing, or importing into other systems.

id: ship_asset_to_file
namespace: kestra.company.data
tasks:
- id: export_assets
type: io.kestra.plugin.ee.assets.AssetShipper
assetExporters:
- id: file_exporter
type: io.kestra.plugin.ee.assets.FileAssetExporter
format: ION

You can change the format property to JSON if you prefer a more widely-compatible format.

Export assets to OpenLineage

Ship asset metadata to an OpenLineage-compatible lineage provider. This requires mapping Kestra asset fields to OpenLineage conventions.

id: ship_asset_to_openlineage
namespace: kestra.company.data
tasks:
- id: export_to_lineage
type: io.kestra.plugin.ee.assets.AssetShipper
assetExporters:
- id: openlineage_exporter
type: io.kestra.plugin.ee.openlineage.OpenLineageAssetExporter
uri: http://host.docker.internal:5000
mappings:
io.kestra.plugin.ee.assets.Table:
namespace: namespace

The mappings property defines how Kestra asset metadata fields map to OpenLineage dataset facets. Each asset type can have its own mapping configuration. For more information about OpenLineage dataset facets and available fields, see the OpenLineage Dataset Facets documentation.

Purge assets and lineage data

The io.kestra.plugin.ee.assets.PurgeAssets task enforces asset retention without touching executions or logs. By default, this task purges assets, asset usage events (execution view), and asset lineage events (for asset exporters) matching the filters. You can configure it to only purge specific types of records.

Filters:

PropertyDescription
namespaceFilter by namespace. Supports prefix matching (e.g., company.data matches company.data.staging).
assetIdFilter by a specific asset ID.
assetTypeFilter by one or more asset types (e.g., io.kestra.plugin.ee.assets.Table).
metadataQueryFilter by metadata key-value pairs.
endDate(required) Purge records created or updated before this date (ISO 8601).

Purge scope:

PropertyDefaultDescription
purgeAssetstrueWhether to purge the asset records themselves.
purgeAssetUsagestrueWhether to purge asset usage events (execution view).
purgeAssetLineagestrueWhether to purge asset lineage events.

Outputs: purgedAssetsCount, purgedAssetUsagesCount, purgedAssetLineagesCount.

Example: purge old VM assets on a monthly schedule.

id: asset_retention_policy
namespace: company.infra
triggers:
- id: monthly_cleanup
type: io.kestra.plugin.core.trigger.Schedule
cron: "0 0 1 * *"
tasks:
- id: purge_old_vms
type: io.kestra.plugin.ee.assets.PurgeAssets
assetType:
- io.kestra.plugin.ee.assets.VM
endDate: "{{ now() | dateAdd(-180, 'DAYS') }}"

Visualizing assets in dashboards

The io.kestra.plugin.ee.dashboard.data.Assets data source builds charts over the asset inventory directly in a custom dashboard. Asset charts are not filtered by the dashboard time range; they always reflect the current state of your inventory.

See Assets (EE and Cloud only) in the Dashboards documentation for available fields, chart type compatibility, and configuration examples.

Was this page helpful?