Skip to content
1 change: 0 additions & 1 deletion morpheus/utils/nvt/patches/merlin_patches.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
104 changes: 54 additions & 50 deletions morpheus/utils/nvt/schema_converters.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand All @@ -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.

Expand Down Expand Up @@ -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.

Expand All @@ -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.

Expand All @@ -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.

Expand All @@ -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.

Expand All @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -319,15 +323,15 @@ 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:
node = queue.pop(0)
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
Expand All @@ -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)

Expand All @@ -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
Expand All @@ -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",
Expand Down
3 changes: 1 addition & 2 deletions morpheus/utils/nvt/transforms.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@
# See the License for the specific language governing permissions and
# limitations under the License.


import json
import typing

Expand All @@ -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):
Expand Down
12 changes: 6 additions & 6 deletions morpheus/utils/schema_transforms.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -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):
Expand Down Expand Up @@ -139,5 +139,5 @@ def process_dataframe(

if (convert_to_pd):
return result.to_pandas()
else:
return result

return result
12 changes: 5 additions & 7 deletions tests/examples/digital_fingerprinting/test_dfp_file_to_df.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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',
Expand All @@ -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)

Expand Down
28 changes: 13 additions & 15 deletions tests/examples/ransomware_detection/test_create_features.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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']
Expand Down Expand Up @@ -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)

Expand Down
Loading