Skip to content

SFTPSensor.newer_than not working with jinja logical ds/ts expression #36629

Description

@fpopic

Apache Airflow Provider(s)

sftp

Versions of Apache Airflow Providers

apache-airflow-providers-sftp==4.4.0

Apache Airflow version

2.6.3

Operating System

Ubuntu 20.04.6 LTS

Deployment

Virtualenv installation

Deployment details

No response

What happened

I tried to use parameter newer_than of type datetime with jinja expression for logical execution date {{ ds }} or timestamp {{ ts }} but couldn't find a single way how to please the SFTPSensor.newer_than checks.

import datetime
import pendulum

from airflow import models
from airflow.providers.sftp.sensors.sftp import SFTPSensor
from airflow.operators.empty import EmptyOperator


with models.DAG(
    "dag_example_sftp_sensor_newer_than_example",
    schedule_interval="@once",
    start_date=datetime.datetime(2024, 1, 1)
) as dag:

    start_dag = EmptyOperator(task_id="start_dag")
    end_dag = EmptyOperator(task_id="end_dag")

    wait_for_sftp_file = SFTPSensor(
        task_id=f"wait_for_sftp_file",
        sftp_conn_id="sftp_conn_id",
        path=f"some-other-jinja-expression-depending-on-airflow-{{{{ ds }}}}",
        newer_than=pendulum.from_format('{{ ds }}', fmt="YYYY-MM-DD")
    )

    start_dag >> wait_for_sftp_file >> end_dag

I get

File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/dagbag.py", line 346, in parse
    loader.exec_module(new_module)
  File "<frozen importlib._bootstrap_external>", line 843, in exec_module
  File "<frozen importlib._bootstrap>", line 219, in _call_with_frames_removed
  File "/home/airflow/gcs/dags/dag_example_sftp_sensor_newer_than_example.py", line XXX, in <module>
    newer_than=pendulum.from_format('{{ ds }}', fmt="YYYY-MM-DD")
  File "/opt/python3.8/lib/python3.8/site-packages/pendulum/__init__.py", line 259, in from_format
    parts = _formatter.parse(string, fmt, now(), locale=locale)
  File "/opt/python3.8/lib/python3.8/site-packages/pendulum/formatting/formatter.py", line 413, in parse
    raise ValueError("String does not match format {}".format(fmt))
ValueError: String does not match format YYYY-MM-DD

In case I hard-code the value (not using jinja) to something like

newer_than=pendulum.from_format('2024-01-01', fmt="YYYY-MM-DD")

everything works.

What you think should happen instead

Parameter newer_than should be working with jinja templates {{ts}} or {{ds}}.

How to reproduce

airflow tasks test dag_example_sftp_sensor_newer_than_example wait_for_sftp_file 2024-01-01

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 Jan 6, 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. changed the title [-]SFTPSensor.newer_than not working with jinja logical date expression[/-] [+][WIP] SFTPSensor.newer_than not working with jinja logical date expression[/+] on Jan 6, 2024
  3. vatsrahul1001 commented on Jan 6, 2024

    @vatsrahul1001
    Contributor

    @fpopic instead of using pendulum.from_format('{{ ds }}', fmt="YYYY-MM-DD") can you try newer_than='{{ ds }}' as defauft format is YYYY-MM-DD only

  4. removed
    needs-triagelabel for new issues that we didn't triage yet
    on Jan 7, 2024
  5. RNHTTR commented on Jan 7, 2024

    @RNHTTR
    Contributor

    For what it's worth, the SFTPSensor does support templating newer_than, so to @vatsrahul1001 point, I wonder if it's something to do with when the DateTime is parsed?

  6. fpopic commented on Jan 7, 2024

    @fpopic
    ContributorAuthor

    @fpopic instead of using pendulum.from_format('{{ ds }}', fmt="YYYY-MM-DD") can you try newer_than='{{ ds }}' as defauft format is YYYY-MM-DD only

    @vatsrahul1001 this works.

    I found a simpler way using pendulum https://pendulum.eustace.io/docs/#addition-and-subtraction newer_than='{{ logical_date.add(days=-1, hours=-2) }}'

  7. changed the title [-][WIP] SFTPSensor.newer_than not working with jinja logical date expression[/-] [+]SFTPSensor.newer_than not working with jinja logical date expression[/+] on Jan 7, 2024
  8. changed the title [-]SFTPSensor.newer_than not working with jinja logical date expression[/-] [+]SFTPSensor.newer_than not working with jinja logical ds/ts expression[/+] on Jan 7, 2024
  9. fpopic commented on Jan 9, 2024

    @fpopic
    ContributorAuthor

    @RNHTTR @vatsrahul1001 Sorry for closing, but was too early.

    My pattern was testing for non-existing file on SFTP server and it did not raise import exception, since the newer_than check comes later in code, but once the file exists on SFTP server, I am getting

    [2024-01-09T08:18:29.763+0000] {sftp.py:90} INFO - Found File XXXX_YYY_KW2351_ZZZZ.WWW last modified: 20231221050445
    [2024-01-09 08:18:29,771] {taskinstance.py:1826} ERROR - Task failed with exception\nTraceback (most recent call last):\n  File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 225, in execute\n    raise e\n  File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 212, in execute\n    poke_return = self.poke(context)\n  File "/home/airflow/.local/lib/python3.8/site-packages/airflow/providers/sftp/sensors/sftp.py", line 97, in poke\n    _newer_than = convert_to_utc(self.newer_than)\n  File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 104, in convert_to_utc\n    if not is_localized(value):\n  File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 39, in is_localized\n    return value.utcoffset() is not None\nAttributeError: 'str' object has no attribute 'utcoffset'
    [2024-01-09T08:18:29.771+0000] {taskinstance.py:1826} ERROR - Task failed with exception
    Traceback (most recent call last):
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 225, in execute
        raise e
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 212, in execute
        poke_return = self.poke(context)
      File "/home/airflow/.local/lib/python3.8/site-packages/airflow/providers/sftp/sensors/sftp.py", line 97, in poke
        _newer_than = convert_to_utc(self.newer_than)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 104, in convert_to_utc
        if not is_localized(value):
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 39, in is_localized
        return value.utcoffset() is not None
    AttributeError: 'str' object has no attribute 'utcoffset'
    [2024-01-09 08:18:29,783] {taskinstance.py:1346} INFO - Marking task as FAILED. dag_id=dag_XXX, task_id=wait_for_sftp_file, execution_date=20231222T000000, start_date=, end_date=20240109T081829
    [2024-01-09T08:18:29.783+0000] {taskinstance.py:1346} INFO - Marking task as FAILED. dag_id=dag_XXX, task_id=wait_for_sftp_file, execution_date=20231222T000000, start_date=, end_date=20240109T081829
    Traceback (most recent call last):
      File "/opt/python3.8/bin/airflow", line 8, in <module>
        sys.exit(main())
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/__main__.py", line 48, in main
        args.func(args)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/cli/cli_config.py", line 52, in command
        return func(*args, **kwargs)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/cli.py", line 112, in wrapper
        return f(*args, **kwargs)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/cli/commands/task_command.py", line 618, in task_test
        ti.run(ignore_task_deps=True, ignore_ti_state=True, test_mode=True)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/session.py", line 76, in wrapper
        return func(*args, session=session, **kwargs)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1723, in run
        self._run_raw_task(
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/session.py", line 73, in wrapper
        return func(*args, **kwargs)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1408, in _run_raw_task
        self._execute_task_with_callbacks(context, test_mode)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1560, in _execute_task_with_callbacks
        result = self._execute_task(context, task_orig)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1630, in _execute_task
        result = execute_callable(context=context)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 225, in execute
        raise e
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/sensors/base.py", line 212, in execute
        poke_return = self.poke(context)
      File "/home/airflow/.local/lib/python3.8/site-packages/airflow/providers/sftp/sensors/sftp.py", line 97, in poke
        _newer_than = convert_to_utc(self.newer_than)
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 104, in convert_to_utc
        if not is_localized(value):
      File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timezone.py", line 39, in is_localized
        return value.utcoffset() is not None
    AttributeError: 'str' object has no attribute 'utcoffset'

    this was with:

    wait_for_sftp_file = SFTPSensor(
        task_id=f"wait_for_sftp_file",
        sftp_conn_id="sftp_conn_id",
        path=f"XXXX_YYY_KW{{{{ macros.ds_format(ds, '%Y-%m-%d', '%y%U') }}}}_ZZZZ.WWW",
        newer_than='{{ dag_run.logical_date.add(days=-1, hours=-2) }}',
    )
  10. 12 remaining items

  11. added a commit that references this issue on Apr 19, 2024
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions