diff --git a/morpheus/utils/nvt/patches/merlin_patches.py b/morpheus/utils/nvt/patches/merlin_patches.py index a9f0216d49..fdb48c34ea 100644 --- a/morpheus/utils/nvt/patches/merlin_patches.py +++ b/morpheus/utils/nvt/patches/merlin_patches.py @@ -24,7 +24,6 @@ def patch_numpy_dtype_registry(): Until this is fixed upstream, with the mappings added to `merlin/dtypes/mappings/numpy.py`, this patch should be used. The function is idempotent, and should be called before any NVT operators are used. - :return: """ from merlin.dtypes import _dtype_registry diff --git a/morpheus/utils/nvt/schema_converters.py b/morpheus/utils/nvt/schema_converters.py index a893909aad..b281b626cc 100644 --- a/morpheus/utils/nvt/schema_converters.py +++ b/morpheus/utils/nvt/schema_converters.py @@ -37,19 +37,21 @@ from morpheus.utils.column_info import StringJoinColumn from morpheus.utils.nvt import MutateOp from morpheus.utils.nvt.transforms import json_flatten +from morpheus.utils.type_aliases import DataFrameType def sync_df_as_pandas(func: typing.Callable) -> typing.Callable: """ - This function serves as a decorator that synchronizes cudf.DataFrame to pandas.DataFrame before applying the function. + This function serves as a decorator that synchronizes cudf.DataFrame to pandas.DataFrame before applying the + function. :param func: The function to apply to the DataFrame :return: The wrapped function """ - def wrapper(df: typing.Union[pd.DataFrame, cudf.DataFrame], **kwargs) -> typing.Union[pd.DataFrame, cudf.DataFrame]: + def wrapper(df: DataFrameType, **kwargs) -> DataFrameType: convert_to_cudf = False - if type(df) == cudf.DataFrame: + if isinstance(df, cudf.DataFrame): convert_to_cudf = True df = df.to_pandas() @@ -70,7 +72,7 @@ class JSONFlattenInfo(ColumnInfo): output_col_names: list -def resolve_json_output_columns(input_schema) -> typing.List[typing.Tuple[str, str]]: +def resolve_json_output_columns(input_schema: DataFrameInputSchema) -> typing.List[typing.Tuple[str, str]]: """ Resolves JSON output columns from an input schema. @@ -99,33 +101,32 @@ def resolve_json_output_columns(input_schema) -> typing.List[typing.Tuple[str, s return output_cols -def get_ci_column_selector(ci): +def get_ci_column_selector(col_info: ColumnInfo): """ Return a column selector based on a ColumnInfo object. - :param ci: The ColumnInfo object + :param col_info: The ColumnInfo object :return: A column selector """ - if (ci.__class__ == ColumnInfo): - return ci.name + if (col_info.__class__ == ColumnInfo): + return col_info.name - elif ci.__class__ in [RenameColumn, BoolColumn, DateTimeColumn, StringJoinColumn, IncrementColumn]: - return ci.input_name + if col_info.__class__ in [RenameColumn, BoolColumn, DateTimeColumn, StringJoinColumn, IncrementColumn]: + return col_info.input_name - elif ci.__class__ == StringCatColumn: - return ci.input_columns + if col_info.__class__ == StringCatColumn: + return col_info.input_columns - elif ci.__class__ == JSONFlattenInfo: - return ci.input_col_names + if col_info.__class__ == JSONFlattenInfo: + return col_info.input_col_names - elif ci.__class__ == CustomColumn: + if col_info.__class__ == CustomColumn: return '*' - else: - raise Exception(f"Unknown ColumnInfo type: {ci.__class__}") + raise ValueError(f"Unknown ColumnInfo type: {col_info.__class__}") -def json_flatten_from_input_schema(json_input_cols, json_output_cols) -> MutateOp: +def json_flatten_from_input_schema(json_input_cols: typing.List[str], json_output_cols: typing.List[str]) -> MutateOp: """ Return a JSON flatten operation from an input schema. @@ -139,8 +140,7 @@ def json_flatten_from_input_schema(json_input_cols, json_output_cols) -> MutateO @sync_df_as_pandas -def string_cat_col(df: typing.Union[pd.DataFrame, cudf.DataFrame], output_column, - sep) -> typing.Union[pd.DataFrame, cudf.DataFrame]: +def string_cat_col(df: DataFrameType, output_column: str, sep: str) -> DataFrameType: """ Concatenate the string representation of all supplied columns in a DataFrame. @@ -154,11 +154,12 @@ def string_cat_col(df: typing.Union[pd.DataFrame, cudf.DataFrame], output_column return pd.DataFrame({output_column: cat_col}) -def nvt_string_cat_col(column_selector: ColumnSelector, - df: typing.Union[pd.DataFrame, cudf.DataFrame], - output_column, - input_columns, - sep: str = ', '): +def nvt_string_cat_col( + column_selector: ColumnSelector, # pylint: disable=unused-argument + df: DataFrameType, + output_column: str, + input_columns: typing.List[str], + sep: str = ', '): """ Concatenates the string representation of the specified columns in a DataFrame. @@ -173,8 +174,7 @@ def nvt_string_cat_col(column_selector: ColumnSelector, @sync_df_as_pandas -def increment_column(df: typing.Union[pd.DataFrame, cudf.DataFrame], output_column, input_column, period: str = 'D') \ - -> typing.Union[pd.DataFrame, cudf.DataFrame]: +def increment_column(df: DataFrameType, output_column: str, input_column: str, period: str = 'D') -> DataFrameType: """ Crete an increment a column in a DataFrame. @@ -191,11 +191,15 @@ def increment_column(df: typing.Union[pd.DataFrame, cudf.DataFrame], output_colu def nvt_increment_column(column_selector: ColumnSelector, - df: typing.Union[pd.DataFrame, cudf.DataFrame], - output_column, - input_column, + df: DataFrameType, + output_column: str, + input_column: str, period: str = 'D'): - return increment_column(column_selector, df, output_column, input_column, period) + return increment_column(column_selector=column_selector, + df=df, + output_column=output_column, + input_column=input_column, + period=period) # Mappings from ColumnInfo types to functions that create the corresponding NVT operator @@ -281,14 +285,14 @@ def build_nx_dependency_graph(column_info_objects: typing.List[ColumnInfo]) -> n graph = nx.DiGraph() def find_dependent_column(name, current_name): - for ci in column_info_objects: - if ci.name == current_name: + for col_info in column_info_objects: + if col_info.name == current_name: continue - if ci.name == name: - return ci - elif ci.__class__ == JSONFlattenInfo: - if name in [c for c, _ in ci.output_col_names]: - return ci + if col_info.name == name: + return col_info + if col_info.__class__ == JSONFlattenInfo: + if name in [c for c, _ in col_info.output_col_names]: + return col_info return None for col_info in column_info_objects: @@ -319,7 +323,7 @@ def find_dependent_column(name, current_name): def bfs_traversal_with_op_map(graph, ci_map, root_nodes): visited = set() - queue = [n for n in root_nodes] + queue = list(root_nodes) node_op_map = {} while queue: @@ -327,7 +331,7 @@ def bfs_traversal_with_op_map(graph, ci_map, root_nodes): if node not in visited: visited.add(node) - parents = [n for n in graph.predecessors(node)] + parents = list(graph.predecessors(node)) if len(parents) == 0: # We need to start an operator chain with a column selector, so root nodes need to prepend a parent # column selection operator @@ -346,14 +350,14 @@ def bfs_traversal_with_op_map(graph, ci_map, root_nodes): # Chain ops together into a compound op node_op = parent_input - for op in ops: - node_op = node_op >> op + for operator in ops: + node_op = node_op >> operator # Set the op for this node to the compound operator node_op_map[node] = node_op # Add our neighbors to the queue - neighbors = [n for n in graph.neighbors(node)] + neighbors = list(graph.neighbors(node)) for neighbor in neighbors: queue.append(neighbor) @@ -362,21 +366,21 @@ def bfs_traversal_with_op_map(graph, ci_map, root_nodes): def coalesce_leaf_nodes(node_op_map, graph, preserve_re): coalesced_workflow = None - for node, op in node_op_map.items(): - neighbors = [n for n in graph.neighbors(node)] + for node, operator in node_op_map.items(): + neighbors = list(graph.neighbors(node)) # Only add the operators for leaf nodes, or those explicitly preserved if len(neighbors) == 0 or (preserve_re and preserve_re.match(node)): if coalesced_workflow is None: - coalesced_workflow = op + coalesced_workflow = operator else: - coalesced_workflow = coalesced_workflow + op + coalesced_workflow = coalesced_workflow + operator return coalesced_workflow def coalesce_ops(graph, ci_map, preserve_re=None): root_nodes = [node for node, in_degree in graph.in_degree() if in_degree == 0] - visited, node_op_map = bfs_traversal_with_op_map(graph, ci_map, root_nodes) + _, node_op_map = bfs_traversal_with_op_map(graph, ci_map, root_nodes) coalesced_workflow = coalesce_leaf_nodes(node_op_map, graph, preserve_re=preserve_re) return coalesced_workflow @@ -400,11 +404,11 @@ def dataframe_input_schema_to_nvt_workflow(input_schema: DataFrameInputSchema, v json_output_cols = resolve_json_output_columns(input_schema) json_cols = input_schema.json_columns - column_info_objects = [ci for ci in input_schema.column_info] + column_info_objects = list(input_schema.column_info) if (json_cols is not None and len(json_cols) > 0): column_info_objects.append( JSONFlattenInfo( - input_col_names=[c for c in json_cols], + input_col_names=list(json_cols), # output_col_names=[name for name, _ in json_output_cols], output_col_names=json_output_cols, dtype="str", diff --git a/morpheus/utils/nvt/transforms.py b/morpheus/utils/nvt/transforms.py index 1f350c2e5c..26de333437 100644 --- a/morpheus/utils/nvt/transforms.py +++ b/morpheus/utils/nvt/transforms.py @@ -12,7 +12,6 @@ # See the License for the specific language governing permissions and # limitations under the License. - import json import typing @@ -32,7 +31,7 @@ def json_flatten(col_selector, df: typing.Union[pd.DataFrame, cudf.DataFrame]): pd_series = df[col] if not convert_to_cudf else df[col].to_pandas() pd_series = pd_series.apply(lambda x: x if isinstance(x, dict) else json.loads(x)) pdf_norm = pd.json_normalize(pd_series) - pdf_norm.rename(columns=lambda x: col + "." + x, inplace=True) + pdf_norm.rename(columns=lambda x, c=col: c + "." + x, inplace=True) pdf_norm.reset_index(drop=True, inplace=True) if (df_normalized is None): diff --git a/morpheus/utils/schema_transforms.py b/morpheus/utils/schema_transforms.py index b060d0d1ee..1b7e2d7548 100644 --- a/morpheus/utils/schema_transforms.py +++ b/morpheus/utils/schema_transforms.py @@ -30,7 +30,7 @@ patch_numpy_dtype_registry() # ======================================================================== -logger = logging.getLogger("morpheus.{}".format(__name__)) +logger = logging.getLogger(f"morpheus.{__name__}") def _process_columns(df_in, input_schema: DataFrameInputSchema): @@ -44,11 +44,11 @@ def _process_columns(df_in, input_schema: DataFrameInputSchema): convert_to_cudf = True # Iterate over the column info - for ci in input_schema.column_info: + for col_info in input_schema.column_info: try: - output_df[ci.name] = ci._process_column(df_in) + output_df[col_info.name] = col_info._process_column(df_in) except Exception: - logger.exception("Failed to process column '%s'. Dataframe: \n%s", ci.name, df_in, exc_info=True) + logger.exception("Failed to process column '%s'. Dataframe: \n%s", col_info.name, df_in, exc_info=True) raise if (input_schema.preserve_columns is not None): @@ -139,5 +139,5 @@ def process_dataframe( if (convert_to_pd): return result.to_pandas() - else: - return result + + return result diff --git a/tests/examples/digital_fingerprinting/test_dfp_file_to_df.py b/tests/examples/digital_fingerprinting/test_dfp_file_to_df.py index 54031fac8a..054616d889 100644 --- a/tests/examples/digital_fingerprinting/test_dfp_file_to_df.py +++ b/tests/examples/digital_fingerprinting/test_dfp_file_to_df.py @@ -22,8 +22,6 @@ import pandas as pd import pytest -import cudf - from morpheus.common import FileTypes from morpheus.config import Config from morpheus.pipeline.preallocator_mixin import PreallocatorMixin @@ -34,8 +32,8 @@ from utils.dataset_manager import DatasetManager -@pytest.fixture -def single_file_obj(): +@pytest.fixture(name='single_file_obj') +def single_file_obj_fixture(): input_file = os.path.join(TEST_DIRS.tests_data_dir, 'appshield', 'snapshot-1', @@ -55,10 +53,10 @@ def test_single_object_to_dataframe(single_file_obj: fsspec.core.OpenFile): assert df.columns == ['data'] with open(single_file_obj.path, encoding='UTF-8') as fh: - d = json.load(fh) - expected_data = d['data'] + json_data = json.load(fh) + expected_data = json_data['data'] - aslist = [x.tolist() for x in df['data'].to_list()] # to_list returns a list of numpy arrays + aslist = [x.tolist() for x in df['data'].to_list()] # to_list returns a list of numpy arrays assert (aslist == expected_data) diff --git a/tests/examples/ransomware_detection/test_create_features.py b/tests/examples/ransomware_detection/test_create_features.py index 9be03d4838..97e3e4063d 100644 --- a/tests/examples/ransomware_detection/test_create_features.py +++ b/tests/examples/ransomware_detection/test_create_features.py @@ -32,14 +32,16 @@ @pytest.mark.use_python class TestCreateFeaturesRWStage: + # pylint: disable=no-name-in-module @mock.patch('stages.create_features.Client') - def test_constructor(self, - mock_dask_client, - config: Config, - dask_distributed: types.ModuleType, - rwd_conf: dict, - interested_plugins: typing.List[str]): + def test_constructor( + self, + mock_dask_client, + config: Config, + dask_distributed: types.ModuleType, # pylint: disable=unused-argument + rwd_conf: dict, + interested_plugins: typing.List[str]): mock_dask_client.return_value = mock_dask_client from common.data_models import FeatureConfig from common.feature_extractor import FeatureExtractor @@ -56,10 +58,6 @@ def test_constructor(self, assert isinstance(stage, MultiMessageStage) assert stage._client is mock_dask_client - scheduler_info = stage._client.scheduler_info() - len(scheduler_info['workers']) == n_workers - for worker in scheduler_info['workers'].values(): - assert worker['nthreads'] == threads_per_worker assert isinstance(stage._feature_config, FeatureConfig) assert stage._feature_config.file_extns == rwd_conf['file_extensions'] @@ -172,12 +170,12 @@ def test_create_multi_messages(self, assert len(multi_messages) == len(pids) prev_loc = 0 - for (i, mm) in enumerate(multi_messages): - assert isinstance(mm, MultiMessage) + for (i, multi_msg) in enumerate(multi_messages): + assert isinstance(multi_msg, MultiMessage) pid = pids[i] - (mm.get_meta(['pid_process']) == pid).all() - assert mm.mess_offset == prev_loc - prev_loc = mm.mess_offset + mm.mess_count + (multi_msg.get_meta(['pid_process']) == pid).all() + assert multi_msg.mess_offset == prev_loc + prev_loc = multi_msg.mess_offset + multi_msg.mess_count assert prev_loc == len(df) diff --git a/tests/test_column_info.py b/tests/test_column_info.py index a50659198b..a576b119d1 100644 --- a/tests/test_column_info.py +++ b/tests/test_column_info.py @@ -33,6 +33,7 @@ from morpheus.utils.column_info import RenameColumn from morpheus.utils.column_info import StringCatColumn from morpheus.utils.column_info import StringJoinColumn +from morpheus.utils.nvt import dataframe_input_schema_to_nvt_workflow from morpheus.utils.schema_transforms import process_dataframe from utils import TEST_DIRS @@ -99,10 +100,9 @@ def test_dataframe_input_schema_with_json_cols(): assert "properties.userPrincipalName" not in processed_df_cols # Test that we get the same answer when the dataframe is processed via workflow - # TODO: Uncomment once https://github.com/rapidsai/cudf/issues/13305 is fixed - # nvt_workflow = dataframe_input_schema_to_nvt_workflow(schema) - # df_processed_workflow = process_dataframe(input_df, nvt_workflow) - # assert df_processed_schema.equals(df_processed_workflow) + nvt_workflow = dataframe_input_schema_to_nvt_workflow(schema) + df_processed_workflow = process_dataframe(input_df, nvt_workflow) + assert df_processed_schema.equals(df_processed_workflow) @pytest.mark.use_python diff --git a/tests/utils/nvt/integration/test_mutate_op.py b/tests/utils/nvt/integration/test_mutate_op.py index 028c143eda..bad0322030 100644 --- a/tests/utils/nvt/integration/test_mutate_op.py +++ b/tests/utils/nvt/integration/test_mutate_op.py @@ -46,9 +46,9 @@ def test_integration_pandas(json_data: typing.List[str], expected_pdf: pd.DataFr pdf = pd.DataFrame({'col1': json_data}) col_selector = ColumnSelector(['col1']) - op = MutateOp(json_flatten, [("col1.key1", "object"), ("col1.key2.subkey1", "object"), - ("col1.key2.subkey2", "object")]) - result_pdf = op.transform(col_selector, pdf) + mutate_op = MutateOp(json_flatten, [("col1.key1", "object"), ("col1.key2.subkey1", "object"), + ("col1.key2.subkey2", "object")]) + result_pdf = mutate_op.transform(col_selector, pdf) assert result_pdf.equals(expected_pdf), "Integration test with pandas DataFrame failed" @@ -57,9 +57,9 @@ def test_integration_cudf(json_data: typing.List[str], expected_pdf: pd.DataFram cdf = cudf.DataFrame({'col1': json_data}) col_selector = ColumnSelector(['col1']) - op = MutateOp(json_flatten, [("col1.key1", "object"), ("col1.key2.subkey1", "object"), - ("col1.key2.subkey2", "object")]) - result_cdf = op.transform(col_selector, cdf) + mutate_op = MutateOp(json_flatten, [("col1.key1", "object"), ("col1.key2.subkey1", "object"), + ("col1.key2.subkey2", "object")]) + result_cdf = mutate_op.transform(col_selector, cdf) result_pdf = result_cdf.to_pandas() assert result_pdf.equals(expected_pdf), "Integration test with cuDF DataFrame failed" diff --git a/tests/utils/nvt/test_mutate_op.py b/tests/utils/nvt/test_mutate_op.py index d891d63b0f..8c6c8bef52 100644 --- a/tests/utils/nvt/test_mutate_op.py +++ b/tests/utils/nvt/test_mutate_op.py @@ -36,9 +36,9 @@ def example_transform(col_selector: ColumnSelector, df: DataFrameType) -> DataFr def test_transform(df: DataFrameType): - op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) + mutate_op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) col_selector = ColumnSelector(['A', 'B']) - transformed_df = op.transform(col_selector, df) + transformed_df = mutate_op.transform(col_selector, df) expected_df = df.copy() expected_df['A_new'] = df['A'] * 2 @@ -53,12 +53,12 @@ def test_transform(df: DataFrameType): # Test for lambda function transformation def test_transform_lambda(df: DataFrameType): - op = MutateOp(lambda col_selector, - df: df.assign(**{f"{col}_new": df[col] * 2 - for col in col_selector.names}), - output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) + mutate_op = MutateOp(lambda col_selector, + df: df.assign(**{f"{col}_new": df[col] * 2 + for col in col_selector.names}), + output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) col_selector = ColumnSelector(['A', 'B']) - transformed_df = op.transform(col_selector, df) + transformed_df = mutate_op.transform(col_selector, df) expected_df = df.copy() expected_df['A_new'] = df['A'] * 2 @@ -76,10 +76,11 @@ def additional_transform(col_selector: ColumnSelector, df: DataFrameType) -> Dat df['D'] = df['A'] + df['B'] return df - op = MutateOp(additional_transform, - output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64')), ('D', np.dtype('int64'))]) + mutate_op = MutateOp(additional_transform, + output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64')), + ('D', np.dtype('int64'))]) col_selector = ColumnSelector(['A', 'B']) - transformed_df = op.transform(col_selector, df) + transformed_df = mutate_op.transform(col_selector, df) expected_df = df.copy() expected_df['A_new'] = df['A'] * 2 @@ -90,9 +91,9 @@ def additional_transform(col_selector: ColumnSelector, df: DataFrameType) -> Dat def test_column_mapping(): - op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) + mutate_op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) col_selector = ColumnSelector(['A', 'B']) - column_mapping = op.column_mapping(col_selector) + column_mapping = mutate_op.column_mapping(col_selector) expected_mapping = {'A_new': ['A', 'B'], 'B_new': ['A', 'B']} @@ -100,7 +101,7 @@ def test_column_mapping(): def test_compute_output_schema(): - op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) + mutate_op = MutateOp(example_transform, output_columns=[('A_new', np.dtype('int64')), ('B_new', np.dtype('int64'))]) col_selector = ColumnSelector(['A', 'B']) input_schema = Schema([ @@ -109,7 +110,7 @@ def test_compute_output_schema(): ColumnSchema('C', dtype=np.dtype('int64')) ]) - output_schema = op.compute_output_schema(input_schema, col_selector) + output_schema = mutate_op.compute_output_schema(input_schema, col_selector) expected_schema = Schema( [ColumnSchema('A_new', dtype=np.dtype('int64')), ColumnSchema('B_new', dtype=np.dtype('int64'))]) diff --git a/tests/utils/nvt/test_schema_converters.py b/tests/utils/nvt/test_schema_converters.py index 932cacc1ac..410389f859 100644 --- a/tests/utils/nvt/test_schema_converters.py +++ b/tests/utils/nvt/test_schema_converters.py @@ -98,76 +98,78 @@ def test_func(df: pd.DataFrame, value: int) -> pd.DataFrame: def test_json_flatten_info_init(): - ci = JSONFlattenInfo(name="json_info", - dtype="str", - input_col_names=["json_col1.a", "json_col2.b"], - output_col_names=["json_output_col1", "json_output_col2"]) - assert ci.name == "json_info" - assert ci.dtype == "str" - assert ci.input_col_names == ["json_col1.a", "json_col2.b"] - assert ci.output_col_names == ["json_output_col1", "json_output_col2"] + col_info = JSONFlattenInfo(name="json_info", + dtype="str", + input_col_names=["json_col1.a", "json_col2.b"], + output_col_names=["json_output_col1", "json_output_col2"]) + assert col_info.name == "json_info" + assert col_info.dtype == "str" + assert col_info.input_col_names == ["json_col1.a", "json_col2.b"] + assert col_info.output_col_names == ["json_output_col1", "json_output_col2"] def test_json_flatten_info_init_missing_input_col_names(): with pytest.raises(TypeError): - ci = JSONFlattenInfo( # noqa F841 + # pylint: disable=no-value-for-parameter + JSONFlattenInfo( # noqa F841 name="json_info", dtype="str", output_col_names=["json_output_col1", "json_output_col2"]) def test_json_flatten_info_init_missing_output_col_names(): with pytest.raises(TypeError): - ci = JSONFlattenInfo(name="json_info", dtype="str", input_col_names=["json_col1.a", "json_col2.b"]) # noqa F841 + # pylint: disable=no-value-for-parameter + JSONFlattenInfo(name="json_info", dtype="str", input_col_names=["json_col1.a", "json_col2.b"]) # noqa F841 def test_get_ci_column_selector_rename_column(): - ci = RenameColumn(input_name="original_name", name="new_name", dtype="str") - result = get_ci_column_selector(ci) + col_info = RenameColumn(input_name="original_name", name="new_name", dtype="str") + result = get_ci_column_selector(col_info) assert result == "original_name" def test_get_ci_column_selector_bool_column(): - ci = BoolColumn(input_name="original_name", - name="new_name", - dtype="bool", - true_values=["True"], - false_values=["False"]) - result = get_ci_column_selector(ci) + col_info = BoolColumn(input_name="original_name", + name="new_name", + dtype="bool", + true_values=["True"], + false_values=["False"]) + result = get_ci_column_selector(col_info) assert result == "original_name" def test_get_ci_column_selector_datetime_column(): - ci = DateTimeColumn(input_name="original_name", name="new_name", dtype="datetime64[ns]") - result = get_ci_column_selector(ci) + col_info = DateTimeColumn(input_name="original_name", name="new_name", dtype="datetime64[ns]") + result = get_ci_column_selector(col_info) assert result == "original_name" def test_get_ci_column_selector_string_join_column(): - ci = StringJoinColumn(input_name="original_name", name="new_name", dtype="str", sep=",") - result = get_ci_column_selector(ci) + col_info = StringJoinColumn(input_name="original_name", name="new_name", dtype="str", sep=",") + result = get_ci_column_selector(col_info) assert result == "original_name" def test_get_ci_column_selector_increment_column(): - ci = IncrementColumn(input_name="original_name", - name="new_name", - dtype="datetime64[ns]", - groupby_column="groupby_col") - result = get_ci_column_selector(ci) + col_info = IncrementColumn(input_name="original_name", + name="new_name", + dtype="datetime64[ns]", + groupby_column="groupby_col") + result = get_ci_column_selector(col_info) assert result == "original_name" def test_get_ci_column_selector_string_cat_column(): - ci = StringCatColumn(name="new_name", dtype="str", input_columns=["col1", "col2"], sep=", ") - result = get_ci_column_selector(ci) + col_info = StringCatColumn(name="new_name", dtype="str", input_columns=["col1", "col2"], sep=", ") + result = get_ci_column_selector(col_info) assert result == ["col1", "col2"] def test_get_ci_column_selector_json_flatten_info(): - ci = JSONFlattenInfo(name="json_info", - dtype="str", - input_col_names=["json_col1.a", "json_col2.b"], - output_col_names=["json_col1_a", "json_col2_b"]) - result = get_ci_column_selector(ci) + col_info = JSONFlattenInfo(name="json_info", + dtype="str", + input_col_names=["json_col1.a", "json_col2.b"], + output_col_names=["json_col1_a", "json_col2_b"]) + result = get_ci_column_selector(col_info) assert result == ["json_col1.a", "json_col2.b"] @@ -199,14 +201,14 @@ def test_resolve_json_output_columns(): def test_resolve_json_output_columns_empty_input_schema(): input_schema = DataFrameInputSchema() output_cols = resolve_json_output_columns(input_schema) - assert output_cols == [] + assert len(output_cols) == 0 def test_resolve_json_output_columns_no_json_columns(): input_schema = DataFrameInputSchema( column_info=[ColumnInfo(name="column1", dtype="int"), ColumnInfo(name="column2", dtype="str")]) output_cols = resolve_json_output_columns(input_schema) - assert output_cols == [] + assert len(output_cols) == 0 def test_resolve_json_output_columns_with_json_columns(): @@ -236,8 +238,8 @@ def test_bfs_traversal_with_op_map(): input_schema = DataFrameInputSchema(json_columns=["access_device", "application", "auth_device", "user"], column_info=source_column_info) - column_info_objects = [ci for ci in input_schema.column_info] - column_info_map = {ci.name: ci for ci in column_info_objects} + column_info_objects = list(input_schema.column_info) + column_info_map = {col_info.name: col_info for col_info in column_info_objects} graph = build_nx_dependency_graph(column_info_objects) root_nodes = [node for node, in_degree in graph.in_degree() if in_degree == 0] visited, node_op_map = bfs_traversal_with_op_map(graph, column_info_map, root_nodes) @@ -253,8 +255,8 @@ def test_coalesce_leaf_nodes(): input_schema = DataFrameInputSchema(json_columns=["access_device", "application", "auth_device", "user"], column_info=source_column_info) - column_info_objects = [ci for ci in input_schema.column_info] - column_info_map = {ci.name: ci for ci in column_info_objects} + column_info_objects = list(input_schema.column_info) + column_info_map = {col_info.name: col_info for col_info in column_info_objects} graph = build_nx_dependency_graph(column_info_objects) root_nodes = [node for node, in_degree in graph.in_degree() if in_degree == 0] @@ -267,7 +269,7 @@ def test_coalesce_leaf_nodes(): # Extract the leaf nodes from the coalesced workflow leaf_nodes = [] - for node, op in node_op_map.items(): + for node in node_op_map: neighbors = list(graph.neighbors(node)) if len(neighbors) == 0: leaf_nodes.append(node) @@ -292,7 +294,7 @@ def test_input_schema_conversion_empty_schema(): empty_schema = DataFrameInputSchema() with pytest.raises(ValueError, match="Input schema is empty"): - workflow = dataframe_input_schema_to_nvt_workflow(empty_schema) # noqa + dataframe_input_schema_to_nvt_workflow(empty_schema) # noqa def test_input_schema_conversion_additional_column(): @@ -534,6 +536,7 @@ def test_input_schema_conversion_with_trivial_filter(): def test_input_schema_conversion_with_functional_filter(): # Create a DataFrameInputSchema instance with the example schema provided + # pylint: disable=singleton-comparison example_schema = DataFrameInputSchema(json_columns=["access_device", "application", "auth_device", "user"], column_info=source_column_info, row_filter=lambda df: df[df["result"] == True]) # noqa E712