From f87923ada083c2a119ab191da040555f2a681d81 Mon Sep 17 00:00:00 2001 From: XD-DENG Date: Mon, 13 Aug 2018 21:55:14 +0800 Subject: [PATCH] [AIRFLOW-2896] Improve HdfsSensor() 1. Make documentation clearer (format for `ignored_ext` should be extensions like ''py rather than '.py'). 2. Ensure upper/lower case would not affect the usage of `ignored_ext` feagure. 3. Add tests for methods filter_for_ignored_ext() and filter_for_filesize(). --- airflow/sensors/hdfs_sensor.py | 7 ++++--- tests/sensors/test_hdfs_sensor.py | 35 +++++++++++++++++++++++++++++++ 2 files changed, 39 insertions(+), 3 deletions(-) diff --git a/airflow/sensors/hdfs_sensor.py b/airflow/sensors/hdfs_sensor.py index 4d95556f47a8d..4175f4cf07849 100644 --- a/airflow/sensors/hdfs_sensor.py +++ b/airflow/sensors/hdfs_sensor.py @@ -81,19 +81,20 @@ def filter_for_ignored_ext(result, ignored_ext, ignore_copying): Will filter if instructed to do so the result to remove matching criteria :param result: (list) of dicts returned by Snakebite ls - :param ignored_ext: (list) of ignored extensions + :param ignored_ext: (list) of ignored extensions, like ``['exe', 'py']`` :param ignore_copying: (bool) shall we ignore ? :return: (list) of dicts which were not removed """ if ignore_copying: log = LoggingMixin().log - regex_builder = "^.*\.(%s$)$" % '$|'.join(ignored_ext) + regex_builder = "^.*\.(%s$)$" % '$|'.join([e.lower() for e in ignored_ext]) ignored_extensions_regex = re.compile(regex_builder) log.debug( 'Filtering result for ignored extensions: %s in files %s', ignored_extensions_regex.pattern, map(lambda x: x['path'], result) ) - result = [x for x in result if not ignored_extensions_regex.match(x['path'])] + result = [x for x in result + if not ignored_extensions_regex.match(x['path'].lower())] log.debug('HdfsSensor.poke: after ext filter result is %s', result) return result diff --git a/tests/sensors/test_hdfs_sensor.py b/tests/sensors/test_hdfs_sensor.py index b94065d8423ef..13b1f3449ca39 100644 --- a/tests/sensors/test_hdfs_sensor.py +++ b/tests/sensors/test_hdfs_sensor.py @@ -89,3 +89,38 @@ def test_legacy_file_does_not_exists(self): # Then with self.assertRaises(AirflowSensorTimeout): task.execute(None) + + def test_filter_for_ignored_ext(self): + """ + Test the method HdfsSensor.filter_for_ignored_ext + :return: + """ + sample_files = [{'path': 'x.py'}, {'path': 'x.txt'}, {'path': 'x.exe'}] + + check_1 = HdfsSensor.filter_for_ignored_ext(result=sample_files, + ignored_ext=['exe', 'py'], + ignore_copying=True) + self.assertTrue(len(check_1) == 1) + self.assertEqual(check_1[0]['path'].rsplit(".")[-1], "txt") + + check_2 = HdfsSensor.filter_for_ignored_ext(result=sample_files, + ignored_ext=['EXE', 'PY'], + ignore_copying=True) + self.assertTrue(len(check_2) == 1) + self.assertEqual(check_2[0]['path'].rsplit(".")[-1], "txt") + + def test_filter_for_filesize(self): + """ + Test the method HdfsSensor.filter_for_filesize + :return: + """ + # unit of 'length' here is "byte" + sample_files = [{'path': 'small_file_1.txt', 'length': 1024}, + {'path': 'small_file_2.txt', 'length': 2048}, + {'path': 'big_file.txt', 'length': 1024 ** 2 + 1}] + + # unit of argument 'size' inside HdfsSensor.filter_for_filesize is "MB" + check = HdfsSensor.filter_for_filesize(result=sample_files, + size=1) + self.assertTrue(len(check) == 1) + self.assertEqual(check[0]['path'], 'big_file.txt')