Zhnz opened a new issue, #73721:
URL: https://github.com/apache/airflow/issues/73721

   ### 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**
   1. Configure Opensearch remote logging
   2. Run a DAG creating over 1000 log messages
   3. Wait for completion
   4. Refresh page to clear live-log and force logs to be pulled from Opensearch
   5. 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:
   
   
https://github.com/apache/airflow/blob/847183dc9f97f1ebcf9016dc1a350c0bd8084843/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py#L1024-L1054
   
   1. 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.
   
   2. 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`.https://github.com/apache/airflow/blob/847183dc9f97f1ebcf9016dc1a350c0bd8084843/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py#L1003
   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?
   
   1. Task logs more than 1k lines
   2. Worker log for this task contains more than 1k lines
   3. Opensearch contains more than 1k lines for this task
   4. After the task is done, the Airflow UI only shows exactly 1k lines, when 
it should show more than 1k lines
   
   ### 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. 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?
   
   - [x] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to