使用Airflow S3 Hook删除S3文件夹内文件时误删目录的问题咨询
解决Airflow S3 Hook误删input/目录的问题
首先得明确一个关键点:S3中并没有真正的“文件夹”——所谓的文件夹只是对象键的前缀。如果你看到控制台里的input/文件夹消失,要么是你删除了键为input/的空对象(手动创建文件夹时会生成这个对象),要么是该前缀下已经没有任何文件对象了,控制台就不再显示这个“文件夹”。
问题根源
你当前的代码在删除后过滤Deleted列表排除input/,但这只是在日志里隐藏了它,删除操作已经执行了。大概率是你的s3_input_keys列表里包含了input/这个键(比如手动创建过文件夹,或者某些操作生成了这个空对象),导致删除请求把它也删掉了。
解决方案:提前过滤要删除的键
核心思路是:在收集要删除的键时,就排除掉input/本身,只保留input/前缀下的文件(即键以input/开头且不等于input/的对象)。
修改后的完整代码示例:
# 1. 先获取input/前缀下的所有键,过滤掉input/本身 s3_input_keys = [ key for key in self._s3_hook.list_keys( bucket=self.s3_bucket, prefix='input/' ) if key != 'input/' ] # 2. 处理无文件可删的情况 if not s3_input_keys: self.log.info("No files to delete in input/ directory") return # 3. 执行删除操作 response_keys = self._s3_hook.delete_objects( bucket=self.s3_bucket, keys=s3_input_keys ) # 4. 处理删除结果 deleted_keys = [x['Key'] for x in response_keys.get("Deleted", [])] self.log.info("Deleted: %s", deleted_keys) if "Errors" in response_keys: errors_keys = [x['Key'] for x in response_keys.get("Errors", [])] raise AirflowException("Errors when deleting: {}".format(errors_keys))
额外说明
- 如果你希望控制台始终显示
input/文件夹(即使没有文件),可以在删除完成后创建一个空的input/对象:
不过这通常没必要,因为当你下次上传文件到self._s3_hook.load_string( string_data='', key='input/', bucket_name=self.s3_bucket )input/前缀下时,控制台会自动重新显示这个“文件夹”。 - 确保你的
list_keys调用正确使用了prefix='input/',这样只会获取该前缀下的对象,不会拿到其他目录的文件。
内容的提问来源于stack exchange,提问作者monster
相关产品推荐
相关产品推荐

