Skip to content

Commit 4da0cc6

Browse files
committed
Refactor: Migrate os.path to pathlib.Path in file_processor_handler.py (#68757)
1 parent a1daaf1 commit 4da0cc6

1 file changed

Lines changed: 21 additions & 22 deletions

File tree

airflow-core/src/airflow/utils/log/file_processor_handler.py

Lines changed: 21 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,6 @@
1818
from __future__ import annotations
1919

2020
import logging
21-
import os
2221
from datetime import datetime
2322
from pathlib import Path
2423

@@ -46,7 +45,7 @@ def __init__(self, base_log_folder, filename_template):
4645
super().__init__()
4746
self.handler = None
4847
self.base_log_folder = base_log_folder
49-
self.dag_dir = os.path.expanduser(settings.DAGS_FOLDER)
48+
self.dag_dir = str(Path(settings.DAGS_FOLDER).expanduser())
5049
self.filename_template, self.filename_jinja_template = parse_template_string(filename_template)
5150

5251
self._cur_date = datetime.today()
@@ -94,9 +93,9 @@ def _render_filename(self, filename):
9493

9594
airflow_directory = airflow.__path__[0]
9695
if filename.startswith(airflow_directory):
97-
filename = os.path.join("native_dags", os.path.relpath(filename, airflow_directory))
96+
filename = str(Path("native_dags") / Path(filename).relative_to(airflow_directory))
9897
else:
99-
filename = os.path.relpath(filename, self.dag_dir)
98+
filename = str(Path(filename).relative_to(self.dag_dir))
10099
ctx = {"filename": filename}
101100

102101
if self.filename_jinja_template:
@@ -105,7 +104,7 @@ def _render_filename(self, filename):
105104
return self.filename_template.format(filename=ctx["filename"])
106105

107106
def _get_log_directory(self):
108-
return os.path.join(self.base_log_folder, timezone.utcnow().strftime("%Y-%m-%d"))
107+
return str(Path(self.base_log_folder) / timezone.utcnow().strftime("%Y-%m-%d"))
109108

110109
def _symlink_latest_log_directory(self):
111110
"""
@@ -115,22 +114,22 @@ def _symlink_latest_log_directory(self):
115114
116115
:return: None
117116
"""
118-
log_directory = self._get_log_directory()
119-
latest_log_directory_path = os.path.join(self.base_log_folder, "latest")
120-
if os.path.isdir(log_directory):
121-
rel_link_target = Path(log_directory).relative_to(Path(latest_log_directory_path).parent)
117+
log_directory = Path(self._get_log_directory())
118+
latest_log_directory_path = Path(self.base_log_folder) / "latest"
119+
if log_directory.is_dir():
120+
rel_link_target = log_directory.relative_to(latest_log_directory_path.parent)
122121
try:
123122
# if symlink exists but is stale, update it
124-
if os.path.islink(latest_log_directory_path):
125-
if os.path.realpath(latest_log_directory_path) != log_directory:
126-
os.unlink(latest_log_directory_path)
127-
os.symlink(rel_link_target, latest_log_directory_path)
128-
elif os.path.isdir(latest_log_directory_path) or os.path.isfile(latest_log_directory_path):
123+
if latest_log_directory_path.is_symlink():
124+
if latest_log_directory_path.resolve() != log_directory.resolve():
125+
latest_log_directory_path.unlink()
126+
latest_log_directory_path.symlink_to(rel_link_target)
127+
elif latest_log_directory_path.is_dir() or latest_log_directory_path.is_file():
129128
logger.warning(
130129
"%s already exists as a dir/file. Skip creating symlink.", latest_log_directory_path
131130
)
132131
else:
133-
os.symlink(rel_link_target, latest_log_directory_path)
132+
latest_log_directory_path.symlink_to(rel_link_target)
134133
except OSError:
135134
logger.warning("OSError while attempting to symlink the latest log directory")
136135

@@ -141,13 +140,13 @@ def _init_file(self, filename):
141140
:param filename: task instance object
142141
:return: relative log path of the given task instance
143142
"""
144-
relative_log_file_path = os.path.join(self._get_log_directory(), self._render_filename(filename))
145-
log_file_path = os.path.abspath(relative_log_file_path)
146-
directory = os.path.dirname(log_file_path)
143+
relative_log_file_path = Path(self._get_log_directory()) / self._render_filename(filename)
144+
log_file_path = relative_log_file_path.absolute()
145+
directory = log_file_path.parent
147146

148-
Path(directory).mkdir(parents=True, exist_ok=True)
147+
directory.mkdir(parents=True, exist_ok=True)
149148

150-
if not os.path.exists(log_file_path):
151-
open(log_file_path, "a").close()
149+
if not log_file_path.exists():
150+
log_file_path.touch()
152151

153-
return log_file_path
152+
return str(log_file_path)

0 commit comments

Comments
 (0)