Under which category would you file this issue?
Providers
Apache Airflow version
3.3.1
What happened and how to reproduce it?
Issue description
My team is currently evaluating migrating from local logs to remote logging using the Opensearch Provider.
We came across an issue with logging more than a thousand messages within a task instance. While the task instance is running, Airflow reads the live logs from the local logs of the corresponding Airflow worker and all log messages are shown in the Airflow UI. After the task instance completes, the log shown in the UI is no longer requested from the worker itself, but from Opensearch, as is expected. The issue here is, that the final task instance logs, pulled from Opensearch and presented in the UI, are cut-off after 1000 messages.
Steps to reproduce
- Configure Opensearch remote logging
- Run a DAG creating over 1000 log messages
- Wait for completion
- Refresh page to clear live-log and force logs to be pulled from Opensearch
- Observe the log only showing the first 1000 messages
Potential cause
First investigation showed, that the full log is available in both the worker and within the Opensearch index. The data is available, so the problem is with requesting it correctly.
Further investigation into the Opensearch Provider package, specifically the os_task_handler.py OpensearchRemoteLogIO._os_read method, shows there are two possible ways an offset can be set:
|
def _os_read(self, log_id: str, offset: int | str, ti: RuntimeTI) -> OpensearchResponse | None: |
|
"""Return the logs matching ``log_id`` in OpenSearch.""" |
|
query: dict[Any, Any] = { |
|
"query": { |
|
"bool": { |
|
"filter": [{"range": {self.offset_field: {"gt": int(offset)}}}], |
|
"must": [{"match_phrase": {"log_id": log_id}}], |
|
} |
|
} |
|
} |
|
index_patterns = self._get_index_patterns(ti) |
|
try: |
|
max_log_line = self.client.count(index=index_patterns, body=query)["count"] |
|
except NotFoundError as e: |
|
self.log.exception("The target index pattern %s does not exist", index_patterns) |
|
raise e |
|
|
|
if max_log_line != 0: |
|
try: |
|
res = self.client.search( |
|
index=index_patterns, |
|
body=query, |
|
sort=[self.offset_field], |
|
size=self.MAX_LINE_PER_PAGE, |
|
from_=self.MAX_LINE_PER_PAGE * self.PAGE, |
|
) |
|
return OpensearchResponse(self, res) |
|
except Exception as err: |
|
self.log.exception("Could not read log with log_id: %s. Exception: %s", log_id, err) |
|
|
|
return None |
-
The Opensearch client search request range is defined by a size of self.MAX_LINE_PER_PAGE, which is always 1000 and an offset of self.MAX_LINE_PER_PAGE * self.PAGE, with self.PAGE being always 0. Both are hard-coded values, which means here always the first 1000 log messages are requested only. This looks like unfinished pagination code to me.
-
Aside from the from_ parameter of the Opensearch client search method, the offset can also be set via the query defined at the beginning of the _os_read method. This offset is passed as parameter into the method. The _os_read method is only called in one place, where the offset parameter is also hard-coded to 0.
|
response = self._os_read(log_id, 0, ti) |
Since both ways the offset can be defined are hard-coded to 0, it means all logs requested will only ever contain the first 0 to 1000 messages at most.
Proposed solution
Within the OpensearchRemoteLogIO._os_read method, since we know the maximum amount of log lines available before requesting any logs, we can finish the half-baked pagination system and iteratively request the logs in MAX_LINE_PER_PAGE (1000 line) steps until all logs are requested.
During the investigation it also occurred to me, that the entire _os_read method exists duplicated and unused inside the OpensearchTaskHandler class within the same file. Since the OpensearchTaskHandler version is two years older I assume it was copy-pasted to the OpensearchRemoteLogIO, its calls redirected and then forgotten to be deleted. Please correct me if I'm wrong in saying that the older version can be deleted for better code clarity.
What you think should happen instead?
After a task instance is done, the Airflow UI only shows up to 1k lines, when it should show more than 1k lines if more are available.
Operating System
Debian 12 (bookworm)
Deployment
Official Apache Airflow Helm Chart
Apache Airflow Provider(s)
opensearch
Versions of Apache Airflow Providers
apache-airflow-providers-opensearch==1.12.1
Official Helm Chart version
Not Applicable
Kubernetes Version
Not Applicable
Helm Chart configuration
Not Applicable
Docker Image customizations
Added the following packages:
petl
openpyxl
requests
chardet>=3.0.2,<6
psycopg2-binary
statsd
apache-airflow-providers-celery
apache-airflow-providers-standard
apache-airflow-providers-postgres
apache-airflow-providers-fab>=3.7.3
apache-airflow-providers-opensearch==1.12.1
asyncpg
python-ldap
authlib
Anything else?
I'd prefer someone else with more insight take on this issue. Maybe @Owen-CH-Leung or @eladkal can help, since they created the OpensearchRemoteLogIO as far as I can tell (#64364). Thank you for you guys' contribution by the way!
I'd be willing to create a PR if no one else is willing to provide a fix.
Here are our configuration values btw:
logging:
remote_logging: "True"
remote_base_log_folder: "opensearch://"
remote_log_conn_id: "opensearch_default"
delete_local_logs: "False"
opensearch:
host: "[redacted]"
port: 443
username: "${AIRFLOW_OPENSEARCH_USERNAME}"
password: "${AIRFLOW_OPENSEARCH_PASSWORD}"
write_to_os: "False"
write_stdout: "True"
json_format: "True"
target_index: "container"
index_patterns: "container"
log_id_template: "{dag_id}-{task_id}-{run_id}-{map_index}-{try_number}"
offset_field: "offset"
opensearch_configs:
use_ssl: "True"
verify_certs: "True"
ca_certs: "/etc/ssl/certs/internal-ca-certificates.crt"
Are you willing to submit PR?
Code of Conduct
Under which category would you file this issue?
Providers
Apache Airflow version
3.3.1
What happened and how to reproduce it?
Issue description
My team is currently evaluating migrating from local logs to remote logging using the Opensearch Provider.
We came across an issue with logging more than a thousand messages within a task instance. While the task instance is running, Airflow reads the live logs from the local logs of the corresponding Airflow worker and all log messages are shown in the Airflow UI. After the task instance completes, the log shown in the UI is no longer requested from the worker itself, but from Opensearch, as is expected. The issue here is, that the final task instance logs, pulled from Opensearch and presented in the UI, are cut-off after 1000 messages.
Steps to reproduce
Potential cause
First investigation showed, that the full log is available in both the worker and within the Opensearch index. The data is available, so the problem is with requesting it correctly.
Further investigation into the Opensearch Provider package, specifically the
os_task_handler.pyOpensearchRemoteLogIO._os_readmethod, shows there are two possible ways an offset can be set:airflow/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py
Lines 1024 to 1054 in 847183d
The Opensearch client search request range is defined by a size of
self.MAX_LINE_PER_PAGE, which is always1000and an offset ofself.MAX_LINE_PER_PAGE * self.PAGE, withself.PAGEbeing always0. Both are hard-coded values, which means here always the first 1000 log messages are requested only. This looks like unfinished pagination code to me.Aside from the
from_parameter of the Opensearch client search method, the offset can also be set via the query defined at the beginning of the_os_readmethod. This offset is passed as parameter into the method. The_os_readmethod is only called in one place, where the offset parameter is also hard-coded to0.airflow/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py
Line 1003 in 847183d
Since both ways the offset can be defined are hard-coded to
0, it means all logs requested will only ever contain the first 0 to 1000 messages at most.Proposed solution
Within the
OpensearchRemoteLogIO._os_readmethod, since we know the maximum amount of log lines available before requesting any logs, we can finish the half-baked pagination system and iteratively request the logs in MAX_LINE_PER_PAGE (1000 line) steps until all logs are requested.During the investigation it also occurred to me, that the entire
_os_readmethod exists duplicated and unused inside theOpensearchTaskHandlerclass within the same file. Since theOpensearchTaskHandlerversion is two years older I assume it was copy-pasted to theOpensearchRemoteLogIO, its calls redirected and then forgotten to be deleted. Please correct me if I'm wrong in saying that the older version can be deleted for better code clarity.What you think should happen instead?
After a task instance is done, the Airflow UI only shows up to 1k lines, when it should show more than 1k lines if more are available.
Operating System
Debian 12 (bookworm)
Deployment
Official Apache Airflow Helm Chart
Apache Airflow Provider(s)
opensearch
Versions of Apache Airflow Providers
apache-airflow-providers-opensearch==1.12.1
Official Helm Chart version
Not Applicable
Kubernetes Version
Not Applicable
Helm Chart configuration
Not Applicable
Docker Image customizations
Added the following packages:
Anything else?
I'd prefer someone else with more insight take on this issue. Maybe @Owen-CH-Leung or @eladkal can help, since they created the
OpensearchRemoteLogIOas far as I can tell (#64364). Thank you for you guys' contribution by the way!I'd be willing to create a PR if no one else is willing to provide a fix.
Here are our configuration values btw:
Are you willing to submit PR?
Code of Conduct