Skip to content

XComObjectStorageBackend returns the S3 path during deserialization instead of the data #39602

Description

@Uture

Apache Airflow version

2.9.1

If "Other Airflow 2 version" selected, which one?

No response

What happened?

After configuring the object storage as XCom backend, the serialization works fine above the specified threshold, but once another task consumes the previously stored XCom, the deserialization doesn't seem to work. Instead of the deserialized data, the path of the object is returned.

What you think should happen instead?

The stored object should be deserialized and returned to the downstream task.

How to reproduce

import pendulum
from airflow.decorators import dag, task


@dag(
    schedule_interval=None,
    catchup=False,
    start_date=pendulum.datetime(2024, 1, 1, tz="utc"),
)
def dag_test():

    @task()
    def producer():
        import random

        return [random.randint(0, 100) for _ in range(10_000)]

    @task()
    def consumer(obj):
        print(obj)

    producer_t = producer()
    producer_t >> consumer(producer_t)


dag_test()

Operating System

apache/airflow:2.9.0-python3.11

Versions of Apache Airflow Providers

No response

Deployment

Docker-Compose

Deployment details

No response

Anything else?

No response

Are you willing to submit PR?

  • Yes I am willing to submit a PR!

Code of Conduct

Activity

  1. boring-cyborg commented on May 14, 2024

    @boring-cyborg

    Thanks for opening your first issue here! Be sure to follow the issue template! If you are willing to raise PR to address this issue please do so, no need to wait for approval.

  2. kaxil commented on May 14, 2024

    @kaxil
    Member

    cc @bolkedebruin I haven't looked at the code, but was expecting that during deserialization this works out of the box.

    https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/xcoms.html#object-storage-xcom-backend

  3. kaxil commented on May 14, 2024

    @kaxil
    Member

    cc @uranusjr

    def deserialize_value(result: XCom) -> Any:
    """Deserializes the value from the database or object storage.
    Compression is inferred from the file extension.
    """
    data = BaseXCom.deserialize_value(result)
    try:
    path = XComObjectStorageBackend._get_full_path(data)
    except (TypeError, ValueError): # Likely value stored directly in the database.
    return data
    try:
    with path.open(mode="rb", compression="infer") as f:
    return json.load(f, cls=XComDecoder)
    except (TypeError, ValueError):
    return data

  4. uranusjr commented on May 14, 2024

    @uranusjr
    Member

    I seem to recall @TJaniF has a very similar issue a while ago. I don’t remember the details but it was some sort of configuration issues for that instance. And if it is indeed not a bug in Airflow logic, we should try to detect the configuration issue and raise it better to the user, instead of returning a wrong value.

  5. TJaniF commented on May 14, 2024

    @TJaniF
    Contributor

    Yes, I had the same issue and as far as I am aware this PR fixes it: #39313 so it should be fixed in 2.9.2 :)

  6. uranusjr commented on May 14, 2024

    @uranusjr
    Member

    Good to hear Astronomer is now also in the time machine bussiness.

  7. bolkedebruin commented on May 14, 2024

    @bolkedebruin
    Contributor

    Well thanks. That was quick :-).

  8. bolkedebruin commented on May 14, 2024

    @bolkedebruin
    Contributor

    Note: @Uture I don't think the fix will do this for past xcom values. You will need to regenerate those.

  9. Uture commented on May 14, 2024

    @Uture
    Author

    Great, thank you all for resolving this so quickly.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions