From 41a3af0d0652fedf032aad77646d3a2f3b020df3 Mon Sep 17 00:00:00 2001 From: Simon Meierhans Date: Mon, 31 Aug 2026 01:02:02 -0700 Subject: [PATCH] Move TimeseriesPadAndCap into padding stage. PiperOrigin-RevId: 973722138 --- dgf/src/api/data.py | 1 + dgf/src/data/padding.py | 13 +- dgf/src/learning/ten_lines/dataset.py | 1 - dgf/src/transform/BUILD | 32 ++ dgf/src/transform/merge.py | 504 ++++++++++++------- dgf/src/transform/merge_test.py | 388 +++++++++++++- dgf/src/transform/timeseries_padding.py | 367 ++++++++++++++ dgf/src/transform/timeseries_padding_test.py | 453 +++++++++++++++++ 8 files changed, 1579 insertions(+), 180 deletions(-) create mode 100644 dgf/src/transform/timeseries_padding.py create mode 100644 dgf/src/transform/timeseries_padding_test.py diff --git a/dgf/src/api/data.py b/dgf/src/api/data.py index b692870..b0cea7c 100644 --- a/dgf/src/api/data.py +++ b/dgf/src/api/data.py @@ -49,6 +49,7 @@ from dgf.src.data.padding import Padding from dgf.src.data.padding import EdgeSetPadding +from dgf.src.data.padding import FeaturePadding from dgf.src.data.padding import NodeSetPadding from dgf.src.data.graph_snapshots_metadata import GraphSnapshotsFormat diff --git a/dgf/src/data/padding.py b/dgf/src/data/padding.py index 805ab09..931e4f9 100644 --- a/dgf/src/data/padding.py +++ b/dgf/src/data/padding.py @@ -15,17 +15,24 @@ """Padding of graphs.""" import dataclasses -from typing import Dict +from typing import Dict, Optional + + +@dataclasses.dataclass +class FeaturePadding: + max_timeseries_len: Optional[int] = None @dataclasses.dataclass class NodeSetPadding: - num_nodes: int + num_nodes: Optional[int] = None + features: Dict[str, FeaturePadding] = dataclasses.field(default_factory=dict) @dataclasses.dataclass class EdgeSetPadding: - num_edges: int + num_edges: Optional[int] = None + features: Dict[str, FeaturePadding] = dataclasses.field(default_factory=dict) @dataclasses.dataclass diff --git a/dgf/src/learning/ten_lines/dataset.py b/dgf/src/learning/ten_lines/dataset.py index 199f2d4..622f7d8 100644 --- a/dgf/src/learning/ten_lines/dataset.py +++ b/dgf/src/learning/ten_lines/dataset.py @@ -24,7 +24,6 @@ from dgf.src.sampling import config as sampling_config_lib from dgf.src.sampling import in_memory_sampler as in_memory_sampler_lib from dgf.src.transform import merge as merge_lib -from dgf.src.transform import timeseries as timeseries_transform from dgf.src.util import temporal as temporal_util from dgf.src.util import util import numpy as np diff --git a/dgf/src/transform/BUILD b/dgf/src/transform/BUILD index d547c2f..7ac246f 100644 --- a/dgf/src/transform/BUILD +++ b/dgf/src/transform/BUILD @@ -28,10 +28,12 @@ py_library( name = "merge", srcs = ["merge.py"], deps = [ + ":timeseries_padding", "//dgf/src/data:in_memory_graph", "//dgf/src/data:padding", "//dgf/src/data:schema", "//dgf/src/data:tf_in_memory_graph", + "//dgf/src/util:temporal", "//dgf/src/util/weak_dep:weak_dep_tensorflow", # numpy dep, ], @@ -117,6 +119,19 @@ py_library( ], ) +py_library( + name = "timeseries_padding", + srcs = ["timeseries_padding.py"], + deps = [ + "//dgf/src/data:in_memory_graph", + "//dgf/src/data:padding", + "//dgf/src/data:schema", + "//dgf/src/io:feature_format", + "//dgf/src/util:temporal", + # numpy dep, + ], +) + py_library( name = "timeseries", srcs = ["timeseries.py"], @@ -154,8 +169,10 @@ py_test( # apache_beam dep, "//dgf/src/data:in_memory_graph", "//dgf/src/data:padding", + "//dgf/src/data:schema", "//dgf/src/data:tf_in_memory_graph", "//dgf/src/util:gen_test_graph", + "//dgf/src/util:temporal", "//dgf/src/util:test_util", # numpy dep, # tensorflow:tensorflow_no_contrib dep, @@ -261,3 +278,18 @@ py_test( # numpy dep, ], ) + +py_test( + name = "timeseries_padding_test", + srcs = ["timeseries_padding_test.py"], + deps = [ + ":timeseries_padding", + # absl/testing:absltest dep, + # absl/testing:parameterized dep, + "//dgf/src/data:in_memory_graph", + "//dgf/src/data:padding", + "//dgf/src/data:schema", + "//dgf/src/util:test_util", + # numpy dep, + ], +) diff --git a/dgf/src/transform/merge.py b/dgf/src/transform/merge.py index 92ed344..c66f28c 100644 --- a/dgf/src/transform/merge.py +++ b/dgf/src/transform/merge.py @@ -22,6 +22,8 @@ from dgf.src.data import padding as padding_lib from dgf.src.data import schema as schema_lib from dgf.src.data import tf_in_memory_graph +from dgf.src.transform import timeseries_padding +from dgf.src.util import temporal as temporal_util from dgf.src.util.weak_dep.weak_dep_tensorflow import tf import numpy as np @@ -30,12 +32,7 @@ class InsufficientPaddingError(ValueError): pass -def merge_graphs( - graphs: List[in_memory_graph.InMemoryGraph], - schema: schema_lib.GraphSchema, - padding: Optional[padding_lib.Padding], - sentinel_offset: bool = True, -) -> Tuple[in_memory_graph.InMemoryGraph, Dict[str, np.ndarray]]: +class GraphMerger: """Merges multiple `InMemoryGraph` instances into a single graph. Graphs are concatenated sequentially. Node and edge indices from each graph @@ -51,184 +48,318 @@ def merge_graphs( If `padding` is provided, but not large enough to encode all the graphs (e.g., the padding has space for 100 nodes but 120 nodes are provided), - returns a `InsufficientPaddingError` exception. + raises an `InsufficientPaddingError` exception. - Args: - graphs: A list of graphs to merge. - schema: Graph schema. + Attributes: + schema: The original input GraphSchema. padding: The padding configuration to apply. sentinel_offset: Where to include the offset of sentinel nodes / edges added for the padding. Only used when `padding` is set. - - Returns: - A tuple containing: - - The merged `InMemoryGraph`. - - A dictionary mapping each node set name to a NumPy array. The array - at `node_set_offsets[node_set_name][i]` contains the starting node index - of the `i`-th graph from the input `graphs` list within the merged - graph's node set named `node_set_name`. If sentinel_offset=True, - assuming the `graphs` contains `n` values, - all the `node_set_offsets[node_set_name]` will contain `n+1` values. - This extra value represent the fake node used to pad edges. If - sentinel_offset=False, the `node_set_offsets[node_set_name]` will - contain `n` values. + schema_cache: Pre-computed timeseries schema cache. """ - if not graphs: - return in_memory_graph.InMemoryGraph({}, {}), {} - - # TODO(gbm): Failure / warning / crop if the padding is not enought. - # TODO(gbm): TF compatible. - # TODO(gbm): Make padding optionnal. - # TODO(gbm): Implement the graph merging in c++ as part of the sampler. - - # Index of the first node of each graph in the merged graph, for each nodeset. - node_set_offsets: Dict[str, np.ndarray] = {} - for node_set_name in schema.node_sets: - num_nodes = [graph.node_sets[node_set_name].num_nodes for graph in graphs] - offsets = np.cumsum(np.array(num_nodes, dtype=np.int32)) - node_set_offsets[node_set_name] = np.insert(offsets, 0, 0) - - # Determine the edgeset offsets. - # Index of the first edge of each graph in the merged graph, for each nodeset. - edge_set_offsets: Dict[str, np.ndarray] = {} - for edge_set_name in schema.edge_sets: - num_edges = [ - graph.edge_sets[edge_set_name].adjacency.shape[1] for graph in graphs - ] - offsets = np.cumsum(np.array(num_edges, dtype=np.int32)) - edge_set_offsets[edge_set_name] = np.insert(offsets, 0, 0) - - # Merge the nodesets + some more nodeset related compute. - node_set_sentinel_idx: Dict[str, int] = {} - merged_node_sets: Dict[str, in_memory_graph.InMemoryNodeSet] = {} - for node_set_name, node_set_schema in schema.node_sets.items(): - merged_features = {} - node_offsets = node_set_offsets[node_set_name] - # Note: node_offsets[-1] is the number of real nodes - num_real_nodes = node_offsets[-1] - num_sentinel_nodes = None - if padding: - num_nodes = padding.node_sets[node_set_name].num_nodes - # Note: node_offsets[-1] is the number of real nodes - num_sentinel_nodes = num_nodes - num_real_nodes - # The sentinel node is the last one. - node_set_sentinel_idx[node_set_name] = num_nodes - 1 - if num_sentinel_nodes < 1: - raise InsufficientPaddingError( - f"Padding for node set '{node_set_name}' is insufficient. Required" - f" at least {num_real_nodes + 1} nodes (including the sentinel" - f" node), but the padder only defines {num_nodes}." + def __init__( + self, + schema: schema_lib.GraphSchema, + padding: Optional[padding_lib.Padding] = None, + sentinel_offset: bool = True, + schema_cache: Optional[temporal_util.TimeseriesSchemaCache] = None, + ): + self.schema = schema + self.padding = padding + self.sentinel_offset = sentinel_offset + + if padding is not None: + unknown_node_sets = set(padding.node_sets) - set(schema.node_sets) + if unknown_node_sets: + raise ValueError( + f"Padding specifies unknown node sets: {sorted(unknown_node_sets)}." + ) + unknown_edge_sets = set(padding.edge_sets) - set(schema.edge_sets) + if unknown_edge_sets: + raise ValueError( + f"Padding specifies unknown edge sets: {sorted(unknown_edge_sets)}." ) - else: - num_nodes = num_real_nodes - - for feature_name, _ in node_set_schema.features.items(): - # Collect the feature values - merged_values = [ - graph.node_sets[node_set_name].features[feature_name] - for graph in graphs - ] - - if padding: - # Create a feature padding item - first_value = merged_values[0] - padding_shape = list(first_value.shape) - padding_shape[0] = num_sentinel_nodes - - if first_value.dtype.type is np.bytes_: - padding_item = np.full(padding_shape, b"", dtype=first_value.dtype) - else: - padding_item = np.zeros(shape=padding_shape, dtype=first_value.dtype) - merged_values.append(padding_item) - - # Merge together the feature of all the graphs (+ the padding). - merged_features[feature_name] = np.concatenate(merged_values, axis=0) - merged_node_sets[node_set_name] = in_memory_graph.InMemoryNodeSet( - features=merged_features, - num_nodes=num_nodes, + self._has_timeseries_padding = ( + timeseries_padding.has_timeseries_padding(padding) + and temporal_util.schema_has_timeseries_features(schema) ) - # Merge the edges - merged_edge_sets: Dict[str, in_memory_graph.InMemoryEdgeSet] = {} - for edge_set_name, edge_set_schema in schema.edge_sets.items(): - merged_features = {} - # Collect the adjacency of all the graph and apply the offset. - merged_adjacency = [] - for graph_idx, graph in enumerate(graphs): - adjacency = graph.edge_sets[edge_set_name].adjacency - adjacency_offset = np.array([ - [node_set_offsets[edge_set_schema.source][graph_idx]], - [node_set_offsets[edge_set_schema.target][graph_idx]], - ]) - merged_adjacency.append(adjacency + adjacency_offset) - - edge_offsets = edge_set_offsets[edge_set_name] - num_real_edges = edge_offsets[-1] - num_padding = 0 - - if padding: - num_edges = padding.edge_sets[edge_set_name].num_edges - num_padding = num_edges - num_real_edges - if num_padding < 0: - raise InsufficientPaddingError( - f"Padding for edge set '{edge_set_name}' is insufficient. " - f"Required at least {num_real_edges} edges, but the padder only " - f"defines {num_edges}." - ) - - # Add padding edges pointing to sentinel nodes. - padding_node_src = padding.node_sets[edge_set_schema.source].num_nodes - 1 - padding_node_trg = padding.node_sets[edge_set_schema.target].num_nodes - 1 - padding_edge = np.array([[padding_node_src], [padding_node_trg]]) - padding_block = np.tile(padding_edge, (1, num_padding)) - merged_adjacency.append(padding_block) + if self._has_timeseries_padding: + assert padding is not None + if schema_cache is None: + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + self._padded_schema = timeseries_padding.pad_timeseries_schema( + schema, + padding=padding, + schema_cache=schema_cache, + ) else: - num_edges = num_real_edges - - for feature_name, _ in edge_set_schema.features.items(): - if ( - graphs - and feature_name not in graphs[0].edge_sets[edge_set_name].features - ): - continue - # Collect the feature values - merged_values = [ - graph.edge_sets[edge_set_name].features[feature_name] + self._padded_schema = schema + + self.schema_cache = schema_cache + + def output_schema(self) -> schema_lib.GraphSchema: + """Returns the schema of the merged graph.""" + return self._padded_schema + + def __call__( + self, + graphs: List[in_memory_graph.InMemoryGraph], + ) -> Tuple[in_memory_graph.InMemoryGraph, Dict[str, np.ndarray]]: + """Merges a list of graphs into a single graph. + + Args: + graphs: A list of graphs to merge. + + Returns: + A tuple containing: + - The merged `InMemoryGraph`. + - A dictionary mapping each node set name to a NumPy array. The array + at `node_set_offsets[node_set_name][i]` contains the starting node + index of the `i`-th graph from the input `graphs` list within the + merged graph's node set named `node_set_name`. If + sentinel_offset=True, assuming the `graphs` contains `n` values, all + the `node_set_offsets[node_set_name]` will contain `n+1` values. This + extra value represents the fake node used to pad edges. If + sentinel_offset=False, the `node_set_offsets[node_set_name]` will + contain `n` values. + """ + if not graphs: + return in_memory_graph.InMemoryGraph({}, {}), {} + + if self._has_timeseries_padding: + assert self.padding is not None + # TODO(simonmeierhans): Avoid additional memory copy for timeseries + # padding. + graphs = [ + timeseries_padding.pad_timeseries_graph( + graph, + self.schema, + padding=self.padding, + schema_cache=self.schema_cache, + ) for graph in graphs ] - if padding: - # Create a feature padding item - first_value = merged_values[0] - padding_shape = list(first_value.shape) - padding_shape[0] = num_padding - - if first_value.dtype.type is np.bytes_: - padding_item = np.full(padding_shape, b"", dtype=first_value.dtype) - else: - padding_item = np.zeros(shape=padding_shape, dtype=first_value.dtype) - merged_values.append(padding_item) + schema = self._padded_schema + padding = self.padding + sentinel_offset = self.sentinel_offset + + # TODO(gbm): Failure / warning / crop if the padding is not enought. + # TODO(gbm): TF compatible. + # TODO(gbm): Make padding optionnal. + # TODO(gbm): Implement the graph merging in c++ as part of the sampler. + + # Index of the first node of each graph in the merged graph, for each + # nodeset. + node_set_offsets: Dict[str, np.ndarray] = {} + for node_set_name in schema.node_sets: + num_nodes = [ + graph.node_sets[node_set_name].num_nodes for graph in graphs + ] + offsets = np.cumsum(np.array(num_nodes, dtype=np.int32)) + node_set_offsets[node_set_name] = np.insert(offsets, 0, 0) + + # Determine the edgeset offsets. + # Index of the first edge of each graph in the merged graph, for each + # nodeset. + edge_set_offsets: Dict[str, np.ndarray] = {} + for edge_set_name in schema.edge_sets: + num_edges = [ + graph.edge_sets[edge_set_name].adjacency.shape[1] for graph in graphs + ] + offsets = np.cumsum(np.array(num_edges, dtype=np.int32)) + edge_set_offsets[edge_set_name] = np.insert(offsets, 0, 0) + + # Merge the nodesets + some more nodeset related compute. + node_set_sentinel_idx: Dict[str, int] = {} + merged_node_sets: Dict[str, in_memory_graph.InMemoryNodeSet] = {} + for node_set_name, node_set_schema in schema.node_sets.items(): + merged_features = {} + node_offsets = node_set_offsets[node_set_name] + # Note: node_offsets[-1] is the number of real nodes + num_real_nodes = node_offsets[-1] + num_sentinel_nodes = 0 + + target_num_nodes = ( + padding.node_sets[node_set_name].num_nodes + if padding and node_set_name in padding.node_sets + else None + ) + + # Pad the nodeset with sentinel nodes if a padding configuration is + # provided. + if target_num_nodes is not None: + num_nodes = target_num_nodes + # Note: node_offsets[-1] is the number of real nodes + num_sentinel_nodes = num_nodes - num_real_nodes + # The sentinel node is the last one. + node_set_sentinel_idx[node_set_name] = num_nodes - 1 + if num_sentinel_nodes < 1: + raise InsufficientPaddingError( + f"Padding for node set '{node_set_name}' is insufficient. Required" + f" at least {num_real_nodes + 1} nodes (including the sentinel" + f" node), but the padder only defines {num_nodes}." + ) + else: + num_nodes = num_real_nodes + + for feature_name, _ in node_set_schema.features.items(): + # Collect the feature values + merged_values = [ + graph.node_sets[node_set_name].features[feature_name] + for graph in graphs + ] + + if target_num_nodes is not None: + # Create a feature padding item + first_value = merged_values[0] + padding_shape = list(first_value.shape) + padding_shape[0] = num_sentinel_nodes + + if first_value.dtype.type is np.bytes_: + padding_item = np.full(padding_shape, b"", dtype=first_value.dtype) + elif first_value.dtype == np.object_: + padding_item = np.empty(padding_shape, dtype=np.object_) + else: + padding_item = np.zeros( + shape=padding_shape, dtype=first_value.dtype + ) + merged_values.append(padding_item) + + # Merge together the feature of all the graphs (+ the padding). + merged_features[feature_name] = np.concatenate(merged_values, axis=0) + + merged_node_sets[node_set_name] = in_memory_graph.InMemoryNodeSet( + features=merged_features, + num_nodes=num_nodes, + ) + + # Merge the edges + merged_edge_sets: Dict[str, in_memory_graph.InMemoryEdgeSet] = {} + for edge_set_name, edge_set_schema in schema.edge_sets.items(): + merged_features = {} + # Collect the adjacency of all the graph and apply the offset. + merged_adjacency = [] + for graph_idx, graph in enumerate(graphs): + adjacency = graph.edge_sets[edge_set_name].adjacency + adjacency_offset = np.array([ + [node_set_offsets[edge_set_schema.source][graph_idx]], + [node_set_offsets[edge_set_schema.target][graph_idx]], + ]) + merged_adjacency.append(adjacency + adjacency_offset) + + edge_offsets = edge_set_offsets[edge_set_name] + num_real_edges = edge_offsets[-1] + num_padding = 0 + + target_num_edges = ( + padding.edge_sets[edge_set_name].num_edges + if padding and edge_set_name in padding.edge_sets + else None + ) + + # Pad the edgeset with sentinel edges if a padding configuration is + # provided. + if target_num_edges is not None: + num_edges = target_num_edges + num_padding = num_edges - num_real_edges + if num_padding < 0: + raise InsufficientPaddingError( + f"Padding for edge set '{edge_set_name}' is insufficient. " + f"Required at least {num_real_edges} edges, but the padder only " + f"defines {num_edges}." + ) + + # Add padding edges pointing to sentinel nodes. + if ( + edge_set_schema.source not in node_set_sentinel_idx + or edge_set_schema.target not in node_set_sentinel_idx + ): + raise ValueError( + f"Padding for edge set '{edge_set_name}' requires sentinel nodes" + f" on both source '{edge_set_schema.source}' and target" + f" '{edge_set_schema.target}', but one or both node sets are not" + " padded with sentinel nodes." + ) + padding_node_src = node_set_sentinel_idx[edge_set_schema.source] + padding_node_trg = node_set_sentinel_idx[edge_set_schema.target] + adj_dtype = graphs[0].edge_sets[edge_set_name].adjacency.dtype + padding_edge = np.array( + [[padding_node_src], [padding_node_trg]], dtype=adj_dtype + ) + padding_block = np.tile(padding_edge, (1, num_padding)) + merged_adjacency.append(padding_block) + + for feature_name, _ in edge_set_schema.features.items(): + # TODO(simonmeierhans): Remove this check once the sampler supports + # edge features. Subgraphs from samplers currently omit edge features. + if feature_name not in graphs[0].edge_sets[edge_set_name].features: + continue + # Collect the feature values + merged_values = [ + graph.edge_sets[edge_set_name].features[feature_name] + for graph in graphs + ] + + if target_num_edges is not None: + # Create a feature padding item + first_value = merged_values[0] + padding_shape = list(first_value.shape) + padding_shape[0] = num_padding + + if first_value.dtype.type is np.bytes_: + padding_item = np.full(padding_shape, b"", dtype=first_value.dtype) + elif first_value.dtype == np.object_: + padding_item = np.empty(padding_shape, dtype=np.object_) + else: + padding_item = np.zeros( + shape=padding_shape, dtype=first_value.dtype + ) + merged_values.append(padding_item) + + # Merge together the feature of all the graphs (+ the padding). + merged_features[feature_name] = np.concatenate(merged_values, axis=0) + + # Merge all the adjacencies + padding. + merged_adjacency = np.concatenate(merged_adjacency, axis=1) + + merged_edge_sets[edge_set_name] = in_memory_graph.InMemoryEdgeSet( + adjacency=merged_adjacency, + features=merged_features, + ) + + if not sentinel_offset: + node_set_offsets = {k: v[:-1] for k, v in node_set_offsets.items()} + + return ( + in_memory_graph.InMemoryGraph(merged_node_sets, merged_edge_sets), + node_set_offsets, + ) - # Merge together the feature of all the graphs (+ the padding). - merged_features[feature_name] = np.concatenate(merged_values, axis=0) - # Merge all the adjacencies + padding. - merged_adjacency = np.concatenate(merged_adjacency, axis=1) +def merge_graph( + graphs: List[in_memory_graph.InMemoryGraph], + schema: schema_lib.GraphSchema, + padding: Optional[padding_lib.Padding] = None, + sentinel_offset: bool = True, + schema_cache: Optional[temporal_util.TimeseriesSchemaCache] = None, +) -> Tuple[in_memory_graph.InMemoryGraph, Dict[str, np.ndarray]]: + """Merges multiple `InMemoryGraph` instances into a single graph. - merged_edge_sets[edge_set_name] = in_memory_graph.InMemoryEdgeSet( - adjacency=merged_adjacency, - features=merged_features, - ) + Temporary wrapper around the `GraphMerger` class. + """ + return GraphMerger( + schema=schema, + padding=padding, + sentinel_offset=sentinel_offset, + schema_cache=schema_cache, + )(graphs) - if not sentinel_offset: - node_set_offsets = {k: v[:-1] for k, v in node_set_offsets.items()} - return ( - in_memory_graph.InMemoryGraph(merged_node_sets, merged_edge_sets), - node_set_offsets, - ) +merge_graphs = merge_graph def create_padding_item(value, num_padding_items): @@ -265,9 +396,9 @@ def pad_graph_tensorflow( schema: schema_lib.GraphSchema, padding: padding_lib.Padding, ) -> tf_in_memory_graph.TFInMemoryGraph: - """Apply padding on TensorFlow in-memory graphs. + """Pads a TFInMemoryGraph with sentinels according to the padding. - This method is similar to calling merge_graphs with a single graph and a + Similar to the padding implemented in merge_graphs with sentinel_offsets and padder, but work on TFInMemoryGraph instead of (Numpy)InMemoryGraphs. @@ -279,11 +410,24 @@ def pad_graph_tensorflow( Returns: The padded `TFInMemoryGraph`. """ + if temporal_util.schema_has_timeseries_features(schema): + raise NotImplementedError( + "pad_graph_tensorflow does not support timeseries padding yet." + ) + + assert padding.node_sets + assert padding.edge_sets merged_node_sets: Dict[str, tf_in_memory_graph.TFInMemoryNodeSet] = {} for node_set_name, node_set_schema in schema.node_sets.items(): merged_features = {} - num_nodes = padding.node_sets[node_set_name].num_nodes + node_set_padding = padding.node_sets[node_set_name] + if node_set_padding.num_nodes is None: + raise ValueError( + f"Node set '{node_set_name}' does not have num_nodes defined in" + " padding." + ) + num_nodes = node_set_padding.num_nodes num_real_nodes = tf.cast(graph.node_sets[node_set_name].num_nodes, tf.int32) num_sentinel_nodes = num_nodes - num_real_nodes @@ -313,7 +457,13 @@ def pad_graph_tensorflow( for edge_set_name, edge_set_schema in schema.edge_sets.items(): adjacency = graph.edge_sets[edge_set_name].adjacency - num_edges = padding.edge_sets[edge_set_name].num_edges + edge_set_padding = padding.edge_sets[edge_set_name] + if edge_set_padding.num_edges is None: + raise ValueError( + f"Edge set '{edge_set_name}' does not have num_edges defined in" + " padding." + ) + num_edges = edge_set_padding.num_edges num_real_edges = tf.shape(adjacency)[1] num_padding = num_edges - num_real_edges @@ -323,8 +473,16 @@ def pad_graph_tensorflow( message=f"Padding for edge set '{edge_set_name}' is insufficient.", ) - padding_node_src = padding.node_sets[edge_set_schema.source].num_nodes - 1 - padding_node_trg = padding.node_sets[edge_set_schema.target].num_nodes - 1 + src_padding = padding.node_sets[edge_set_schema.source] + trg_padding = padding.node_sets[edge_set_schema.target] + if src_padding.num_nodes is None or trg_padding.num_nodes is None: + raise ValueError( + f"Padding for edge set '{edge_set_name}' requires sentinel nodes on" + f" both source '{edge_set_schema.source}' and target" + f" '{edge_set_schema.target}'." + ) + padding_node_src = src_padding.num_nodes - 1 + padding_node_trg = trg_padding.num_nodes - 1 padding_edge = tf.constant( [[padding_node_src], [padding_node_trg]], dtype=adjacency.dtype @@ -336,6 +494,8 @@ def pad_graph_tensorflow( merged_features = {} for feature_name, _ in edge_set_schema.features.items(): + # TODO(simonmeierhans): Remove this check once the sampler supports + # edge features. Subgraphs from samplers currently omit edge features. if feature_name not in graph.edge_sets[edge_set_name].features: continue value = graph.edge_sets[edge_set_name].features[feature_name] @@ -416,6 +576,8 @@ def remove_padding_sentinels( unpadded_adjacency = adjacency[:, :num_real_edges] unpadded_features = {} for feature_name, _ in edge_set_schema.features.items(): + # TODO(simonmeierhans): Remove this check once the sampler supports + # edge features. Subgraphs from samplers currently omit edge features. if feature_name not in edge_set.features: continue feature_value = edge_set.features[feature_name] diff --git a/dgf/src/transform/merge_test.py b/dgf/src/transform/merge_test.py index 3d45fec..9ad2a5c 100644 --- a/dgf/src/transform/merge_test.py +++ b/dgf/src/transform/merge_test.py @@ -14,12 +14,15 @@ """Test converting between heterogeneous graph edge types.""" +import dataclasses from absl.testing import absltest from dgf.src.data import in_memory_graph as in_memory_graph_lib from dgf.src.data import padding as padding_lib +from dgf.src.data import schema as schema_lib from dgf.src.data import tf_in_memory_graph from dgf.src.transform import merge as merge_lib from dgf.src.util import gen_test_graph +from dgf.src.util import temporal as temporal_util from dgf.src.util import test_util import numpy as np import tensorflow as tf @@ -43,7 +46,9 @@ def test_batch_with_padding(self): "e2": padding_lib.EdgeSetPadding(num_edges=6), }, ) - merged_graph, offsets = merge_lib.merge_graphs(graphs, schema, padding) + merged_graph, offsets = merge_lib.merge_graphs( + graphs, schema, padding=padding + ) expected_merged_graph = in_memory_graph_lib.InMemoryGraph( node_sets={ "n1": in_memory_graph_lib.InMemoryNodeSet( @@ -177,7 +182,7 @@ def test_batch_padding_too_small(self): r"Required at least 5 nodes \(including the sentinel node\), but the" r" padder only defines 4.", ): - _ = merge_lib.merge_graphs(graphs, schema, padding) + _ = merge_lib.merge_graphs(graphs, schema, padding=padding) def test_batch_with_padding_no_sentinel_offset(self): graphs = [ @@ -196,10 +201,10 @@ def test_batch_with_padding_no_sentinel_offset(self): }, ) merged_graph, offsets = merge_lib.merge_graphs( - graphs, schema, padding, sentinel_offset=False + graphs, schema, sentinel_offset=False, padding=padding ) expected_merged_graph, _ = merge_lib.merge_graphs( - graphs, schema, padding, sentinel_offset=True + graphs, schema, sentinel_offset=True, padding=padding ) test_util.assert_are_equal(self, merged_graph, expected_merged_graph) test_util.assert_are_equal( @@ -323,6 +328,34 @@ def pad(graph): " shape.", ) + def test_pad_graph_tensorflow_timeseries_unsupported(self): + tf_graph = gen_test_graph.generate_tf_in_memory_graph( + variable_length=False, + tensor_type="DENSE", + num_nodes_as_tensor=True, + ) + schema = gen_test_graph.generate_schema( + False, False, variable_length=False + ) + schema.node_sets["n1"].features["f1"] = dataclasses.replace( + schema.node_sets["n1"].features["f1"], is_timeseries=True + ) + padding = padding_lib.Padding( + node_sets={ + "n1": padding_lib.NodeSetPadding( + num_nodes=5, + features={ + "f1": padding_lib.FeaturePadding(max_timeseries_len=10) + }, + ) + }, + edge_sets={}, + ) + with self.assertRaisesRegex( + NotImplementedError, "does not support timeseries padding" + ): + merge_lib.pad_graph_tensorflow(tf_graph, schema, padding) + def test_remove_padding_sentinels_with_padding(self): graphs = [ gen_test_graph.generate_in_memory_graph(False, False), @@ -339,7 +372,9 @@ def test_remove_padding_sentinels_with_padding(self): "e2": padding_lib.EdgeSetPadding(num_edges=6), }, ) - merged_graph, offsets = merge_lib.merge_graphs(graphs, schema, padding) + merged_graph, offsets = merge_lib.merge_graphs( + graphs, schema, padding=padding + ) unpadded_graph = merge_lib.remove_padding_sentinels( merged_graph, schema, offsets ) @@ -349,6 +384,349 @@ def test_remove_padding_sentinels_with_padding(self): ) test_util.assert_are_equal(self, unpadded_graph, expected_unpadded_graph) + def test_batch_with_timeseries_padding(self): + schema = schema_lib.GraphSchema( + node_sets={ + "n1": schema_lib.NodeSchema( + features={ + "ts": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.FLOAT_32, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + is_timeseries=True, + shape=(None,), + ), + } + ), + }, + edge_sets={}, + ) + g1 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=2, + features={ + "ts": np.array( + [np.array([1.0, 2.0]), np.array([3.0, 4.0, 5.0])], + dtype=object, + ), + }, + ), + }, + edge_sets={}, + ) + g2 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=1, + features={ + "ts": np.array( + [np.array([6.0, 7.0, 8.0, 9.0, 10.0])], dtype=object + ), + }, + ), + }, + edge_sets={}, + ) + padding = padding_lib.Padding( + node_sets={ + "n1": padding_lib.NodeSetPadding( + num_nodes=4 + 1, + features={ + "ts": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + # Check that schema cache is no longer required and can be omitted. + merged_graph_no_cache, _ = merge_lib.merge_graphs( + [g1, g2], schema, padding=padding, schema_cache=None + ) + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + merged_graph, _ = merge_lib.merge_graphs( + [g1, g2], + schema, + padding=padding, + schema_cache=schema_cache, + ) + test_util.assert_are_equal(self, merged_graph_no_cache, merged_graph) + # Total real nodes: 2 + 1 = 3. Sentinel nodes: 5 - 3 = 2. Total nodes: 5. + self.assertEqual(merged_graph.node_sets["n1"].num_nodes, 5) + ts_feat = merged_graph.node_sets["n1"].features["ts"] + mask_feat = merged_graph.node_sets["n1"].features["ts_mask"] + self.assertEqual(ts_feat.shape, (5, 3)) + self.assertEqual(mask_feat.shape, (5, 3)) + # g1 node 0: [1.0, 2.0] -> padded to [0.0, 1.0, 2.0], mask: [F, T, T] + np.testing.assert_array_equal(ts_feat[0], np.array([0.0, 1.0, 2.0])) + np.testing.assert_array_equal(mask_feat[0], np.array([False, True, True])) + # g1 node 1: [3.0, 4.0, 5.0] -> [3.0, 4.0, 5.0], mask: [T, T, T] + np.testing.assert_array_equal(ts_feat[1], np.array([3.0, 4.0, 5.0])) + np.testing.assert_array_equal(mask_feat[1], np.array([True, True, True])) + # g2 node 0: [6.0, 7.0, 8.0, 9.0, 10.0] -> capped to [8.0, 9.0, 10.0] + np.testing.assert_array_equal(ts_feat[2], np.array([8.0, 9.0, 10.0])) + np.testing.assert_array_equal(mask_feat[2], np.array([True, True, True])) + # Sentinel nodes: [0.0, 0.0, 0.0] + np.testing.assert_array_equal(ts_feat[3], np.array([0.0, 0.0, 0.0])) + np.testing.assert_array_equal(ts_feat[4], np.array([0.0, 0.0, 0.0])) + + def test_batch_with_edge_set_timeseries_padding(self): + schema = schema_lib.GraphSchema( + node_sets={ + "n1": schema_lib.NodeSchema(features={}), + }, + edge_sets={ + "e1": schema_lib.EdgeSchema( + source="n1", + target="n1", + features={ + "weight_ts": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.FLOAT_32, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + is_timeseries=True, + shape=(None,), + ), + }, + ), + }, + ) + g1 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet(num_nodes=2, features={}), + }, + edge_sets={ + "e1": in_memory_graph_lib.InMemoryEdgeSet( + adjacency=np.array([[0], [1]], dtype=np.int32), + features={ + "weight_ts": np.array( + [np.array([1.0, 2.0])], dtype=object + ), + }, + ), + }, + ) + g2 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet(num_nodes=2, features={}), + }, + edge_sets={ + "e1": in_memory_graph_lib.InMemoryEdgeSet( + adjacency=np.array([[0], [1]], dtype=np.int32), + features={ + "weight_ts": np.array( + [np.array([3.0, 4.0, 5.0, 6.0])], dtype=object + ), + }, + ), + }, + ) + padding = padding_lib.Padding( + node_sets={ + "n1": padding_lib.NodeSetPadding(num_nodes=4 + 1), + }, + edge_sets={ + "e1": padding_lib.EdgeSetPadding( + num_edges=3, + features={ + "weight_ts": padding_lib.FeaturePadding( + max_timeseries_len=3 + ) + }, + ), + }, + ) + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + merged_graph, _ = merge_lib.merge_graphs( + [g1, g2], + schema, + padding=padding, + schema_cache=schema_cache, + ) + # Total real edges: 1 + 1 = 2. Total edges with padding: 3. + edge_set = merged_graph.edge_sets["e1"] + self.assertEqual(edge_set.adjacency.shape, (2, 3)) + ts_feat = edge_set.features["weight_ts"] + mask_feat = edge_set.features["weight_ts_mask"] + self.assertEqual(ts_feat.shape, (3, 3)) + self.assertEqual(mask_feat.shape, (3, 3)) + # Edge 0 (g1): [1.0, 2.0] -> [0.0, 1.0, 2.0] + np.testing.assert_array_equal(ts_feat[0], np.array([0.0, 1.0, 2.0])) + np.testing.assert_array_equal(mask_feat[0], np.array([False, True, True])) + # Edge 1 (g2): [3.0, 4.0, 5.0, 6.0] -> [4.0, 5.0, 6.0] + np.testing.assert_array_equal(ts_feat[1], np.array([4.0, 5.0, 6.0])) + np.testing.assert_array_equal(mask_feat[1], np.array([True, True, True])) + # Padding edge (sentinel): [0.0, 0.0, 0.0] + np.testing.assert_array_equal(ts_feat[2], np.array([0.0, 0.0, 0.0])) + + def test_batch_with_timeseries_only_padding(self): + schema = schema_lib.GraphSchema( + node_sets={ + "n1": schema_lib.NodeSchema( + features={ + "ts": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.FLOAT_32, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + is_timeseries=True, + shape=(None,), + ), + } + ), + }, + edge_sets={}, + ) + g1 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=2, + features={ + "ts": np.array( + [np.array([1.0, 2.0]), np.array([3.0, 4.0, 5.0])], + dtype=object, + ), + }, + ), + }, + edge_sets={}, + ) + g2 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=1, + features={ + "ts": np.array( + [np.array([6.0, 7.0, 8.0, 9.0, 10.0])], dtype=object + ), + }, + ), + }, + edge_sets={}, + ) + # num_nodes is None: node topology is unpadded, only timeseries is. + padding = padding_lib.Padding( + node_sets={ + "n1": padding_lib.NodeSetPadding( + features={ + "ts": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + merged_graph, offsets = merge_lib.merge_graphs( + [g1, g2], + schema, + padding=padding, + schema_cache=schema_cache, + ) + # Total real nodes: 2 + 1 = 3. No sentinel nodes added. + self.assertEqual(merged_graph.node_sets["n1"].num_nodes, 3) + ts_feat = merged_graph.node_sets["n1"].features["ts"] + mask_feat = merged_graph.node_sets["n1"].features["ts_mask"] + self.assertEqual(ts_feat.shape, (3, 3)) + self.assertEqual(mask_feat.shape, (3, 3)) + np.testing.assert_array_equal(offsets["n1"], np.array([0, 2, 3])) + # g1 node 0: [1.0, 2.0] -> padded to [0.0, 1.0, 2.0], mask: [F, T, T] + np.testing.assert_array_equal(ts_feat[0], np.array([0.0, 1.0, 2.0])) + np.testing.assert_array_equal(mask_feat[0], np.array([False, True, True])) + # g1 node 1: [3.0, 4.0, 5.0] -> [3.0, 4.0, 5.0], mask: [T, T, T] + np.testing.assert_array_equal(ts_feat[1], np.array([3.0, 4.0, 5.0])) + np.testing.assert_array_equal(mask_feat[1], np.array([True, True, True])) + # g2 node 0: [6.0, 7.0, 8.0, 9.0, 10.0] -> capped to [8.0, 9.0, 10.0] + np.testing.assert_array_equal(ts_feat[2], np.array([8.0, 9.0, 10.0])) + np.testing.assert_array_equal(mask_feat[2], np.array([True, True, True])) + + def test_graph_merger_class_and_output_schema(self): + schema = schema_lib.GraphSchema( + node_sets={ + "n1": schema_lib.NodeSchema( + features={ + "ts": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.FLOAT_32, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + is_timeseries=True, + shape=(None,), + ), + } + ), + }, + edge_sets={}, + ) + g1 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=2, + features={ + "ts": np.array( + [np.array([1.0, 2.0]), np.array([3.0, 4.0, 5.0])], + dtype=np.object_, + ), + }, + ), + }, + edge_sets={}, + ) + g2 = in_memory_graph_lib.InMemoryGraph( + node_sets={ + "n1": in_memory_graph_lib.InMemoryNodeSet( + num_nodes=1, + features={ + "ts": np.array( + [np.array([6.0, 7.0, 8.0, 9.0, 10.0])], + dtype=np.object_, + ), + }, + ), + }, + edge_sets={}, + ) + padding = padding_lib.Padding( + node_sets={ + "n1": padding_lib.NodeSetPadding( + num_nodes=5, + features={ + "ts": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + merger = merge_lib.GraphMerger(schema=schema, padding=padding) + output_schema = merger.output_schema() + self.assertIn("ts", output_schema.node_sets["n1"].features) + self.assertIn("ts_mask", output_schema.node_sets["n1"].features) + self.assertEqual( + output_schema.node_sets["n1"].features["ts"].shape, (3,) + ) + + merged_graph, offsets = merger([g1, g2]) + self.assertEqual(merged_graph.node_sets["n1"].num_nodes, 5) + + # Test merge_graph function alias + merged_graph_alias, offsets_alias = merge_lib.merge_graph( + [g1, g2], schema=schema, padding=padding + ) + test_util.assert_are_equal(self, merged_graph, merged_graph_alias) + test_util.assert_are_equal(self, offsets, offsets_alias) + + def test_unknown_padding_keys_raise_value_error(self): + schema = schema_lib.GraphSchema( + node_sets={"n1": schema_lib.NodeSchema(features={})}, + edge_sets={}, + ) + bad_node_padding = padding_lib.Padding( + node_sets={"unknown_node_set": padding_lib.NodeSetPadding(num_nodes=5)}, + edge_sets={}, + ) + with self.assertRaisesRegex(ValueError, "unknown node sets"): + merge_lib.GraphMerger(schema=schema, padding=bad_node_padding) + + bad_edge_padding = padding_lib.Padding( + node_sets={"n1": padding_lib.NodeSetPadding(num_nodes=5)}, + edge_sets={"unknown_edge_set": padding_lib.EdgeSetPadding(num_edges=5)}, + ) + with self.assertRaisesRegex(ValueError, "unknown edge sets"): + merge_lib.GraphMerger(schema=schema, padding=bad_edge_padding) + if __name__ == "__main__": absltest.main() diff --git a/dgf/src/transform/timeseries_padding.py b/dgf/src/transform/timeseries_padding.py new file mode 100644 index 0000000..9321ca9 --- /dev/null +++ b/dgf/src/transform/timeseries_padding.py @@ -0,0 +1,367 @@ +# Copyright 2022 Google LLC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Padding and capping for timeseries sequence features in graphs.""" + +# pytype: disable=module-attr +import dataclasses +from typing import Any, Dict, List, Optional, Tuple + +from dgf.src.data import in_memory_graph +from dgf.src.data import padding as padding_lib +from dgf.src.data import schema as schema_lib +from dgf.src.io import feature_format +from dgf.src.util import temporal as temporal_util +import numpy as np + + +def _pad_and_cap_single_feature( + raw_series: np.ndarray, + seq_len: int, + feat_shape: Tuple[int, ...], + padding_value: Any, + dtype: Any, + is_static_shape: bool, +) -> Tuple[np.ndarray, np.ndarray]: + """Pads and caps a single sequence feature into (padded_matrix, mask_matrix).""" + num_entities = raw_series.shape[0] + if dtype == np.bytes_: + if raw_series.dtype.kind in ("S", "a"): + dtype = raw_series.dtype + elif raw_series.dtype.kind == "O" and num_entities > 0: + for elem in raw_series: + if isinstance(elem, np.ndarray) and elem.dtype.kind in ("S", "a"): + if dtype == np.bytes_ or elem.dtype.itemsize > dtype.itemsize: + dtype = elem.dtype + if num_entities == 0: + return ( + np.empty((0, seq_len) + feat_shape, dtype=dtype), + np.empty((0, seq_len), dtype=np.bool_), + ) + + # Fast vectorized path when all entities share a fixed sequence length. + if is_static_shape and raw_series.ndim >= 2: + num_steps = raw_series.shape[1] + if num_steps >= seq_len: + padded_matrix = raw_series[:, -seq_len:].astype(dtype, copy=True) + mask_matrix = np.ones((num_entities, seq_len), dtype=np.bool_) + return padded_matrix, mask_matrix + + pad_width = [(0, 0), (seq_len - num_steps, 0)] + [(0, 0)] * len(feat_shape) + padded_matrix = np.pad( + raw_series.astype(dtype, copy=False), + pad_width=pad_width, + mode="constant", + constant_values=padding_value, + ) + mask_width = [(0, 0), (seq_len - num_steps, 0)] + mask_matrix = np.pad( + np.ones((num_entities, num_steps), dtype=np.bool_), + pad_width=mask_width, + mode="constant", + constant_values=False, + ) + return padded_matrix, mask_matrix + + padded_matrix = np.full( + (num_entities, seq_len) + feat_shape, + fill_value=padding_value, + dtype=dtype, + ) + # Binary mask matrix matching sequence length shape (num_entities, seq_len) + # where True indicates valid observed time steps and False indicates + # left-padded steps. + mask_matrix = np.zeros((num_entities, seq_len), dtype=np.bool_) + + # TODO(mesimon): Move into C++ for performance. + for idx in range(num_entities): + raw_arr = raw_series[idx] + if not isinstance(raw_arr, np.ndarray): + raw_arr = np.asarray(raw_arr) + + num_steps = len(raw_arr) + if num_steps >= seq_len: + padded_matrix[idx] = raw_arr[-seq_len:] + mask_matrix[idx] = True + elif num_steps > 0: + padded_matrix[idx, -num_steps:] = raw_arr + mask_matrix[idx, -num_steps:] = True + + return padded_matrix, mask_matrix + + +def has_timeseries_padding( + padding: Optional[padding_lib.Padding], +) -> bool: + """Returns True if padding contains any timeseries feature padding.""" + if padding is None: + return False + for ns in padding.node_sets.values(): + for fp in ns.features.values(): + if fp.max_timeseries_len is not None: + return True + for es in padding.edge_sets.values(): + for fp in es.features.values(): + if fp.max_timeseries_len is not None: + return True + return False + + +def _validate_group_sequence_lengths( + group_specs: List[temporal_util.TimeseriesGroupSpec], + feature_padding: Dict[str, padding_lib.FeaturePadding], + schemas: schema_lib.FeatureSetSchema, +) -> None: + """Validates that all features in each timeseries group share the same sequence length.""" + for group in group_specs: + ts_group = schemas[group.feature_names[0]].group or group.feature_names[0] + padded_features = { + feat: feature_padding[feat].max_timeseries_len + for feat in group.feature_names + if feat in feature_padding + and feature_padding[feat].max_timeseries_len is not None + } + if 0 < len(padded_features) < len(group.feature_names): + missing = [f for f in group.feature_names if f not in padded_features] + raise ValueError( + f"Features in sequence group '{ts_group}' have inconsistent" + f" padding configuration. Missing padding for features: {missing}." + " All features in the same group must share the same sequence length." + ) + if len(padded_features) == len(group.feature_names): + lengths = set(padded_features.values()) + if len(lengths) > 1: + raise ValueError( + f"Features in sequence group '{ts_group}' have conflicting" + f" max_timeseries_len: {padded_features}. All features in the same" + " group must share the same sequence length." + ) + + +def _pad_timeseries_feature_set_schema( + schemas: schema_lib.FeatureSetSchema, + group_specs: List[temporal_util.TimeseriesGroupSpec], + feature_padding: Dict[str, padding_lib.FeaturePadding], +) -> schema_lib.FeatureSetSchema: + """Computes schema for a feature set after padding/capping.""" + _validate_group_sequence_lengths(group_specs, feature_padding, schemas) + new_schemas: schema_lib.FeatureSetSchema = {} + + padded_features = { + feat: fp.max_timeseries_len + for feat, fp in feature_padding.items() + if fp.max_timeseries_len is not None + } + + for feature_name, feature_schema in schemas.items(): + # Skip features that do not have a timeseries sequence length set. + if feature_name not in padded_features: + new_schemas[feature_name] = feature_schema + continue + + seq_len = padded_features[feature_name] + ts_group = feature_schema.group or feature_name + new_schemas[feature_name] = dataclasses.replace( + temporal_util.with_sequence_length(feature_schema, seq_len), + group=ts_group, + ) + # Do not add masks for masks inception. + if feature_schema.semantic == schema_lib.FeatureSemantic.MASK: + continue + + mask_name = temporal_util.get_mask_feature_name(feature_name, schemas) + if mask_name is None: + mask_name = f"{ts_group}_mask" + if mask_name in schemas: + raise ValueError( + f"Cannot generate mask for sequence group '{ts_group}'. The" + f" fallback mask name '{mask_name}' clashes with an existing" + " feature in the schema that is not a valid mask. Please" + " explicitly define a mask feature for this group or rename the" + " clashing feature." + ) + + if mask_name not in new_schemas: + new_schemas[mask_name] = schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.BOOL, + semantic=schema_lib.FeatureSemantic.MASK, + shape=(seq_len,), + is_timeseries=feature_schema.is_timeseries, + group=ts_group, + ) + + return new_schemas + + +def pad_timeseries_schema( + schema: schema_lib.GraphSchema, + padding: padding_lib.Padding, + schema_cache: Optional[temporal_util.TimeseriesSchemaCache] = None, +) -> schema_lib.GraphSchema: + """Returns the GraphSchema after padding/capping timeseries and adding mask features.""" + if not temporal_util.schema_has_timeseries_features(schema): + return schema + + if schema_cache is None: + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + + new_node_sets = {} + for ns_name, ns_schema in schema.node_sets.items(): + ts_specs = schema_cache.node_sets[ns_name] + if not ts_specs: + new_node_sets[ns_name] = ns_schema + else: + new_node_sets[ns_name] = schema_lib.NodeSchema( + features=_pad_timeseries_feature_set_schema( + schemas=ns_schema.features, + group_specs=ts_specs, + feature_padding=padding.node_sets[ns_name].features, + ) + ) + + new_edge_sets = {} + for es_name, es_schema in schema.edge_sets.items(): + ts_specs = schema_cache.edge_sets[es_name] + if not ts_specs: + new_edge_sets[es_name] = es_schema + else: + new_edge_sets[es_name] = schema_lib.EdgeSchema( + source=es_schema.source, + target=es_schema.target, + features=_pad_timeseries_feature_set_schema( + schemas=es_schema.features, + group_specs=ts_specs, + feature_padding=padding.edge_sets[es_name].features, + ), + ) + + return schema_lib.GraphSchema( + node_sets=new_node_sets, edge_sets=new_edge_sets + ) + + +def pad_timeseries_feature_set( + values: in_memory_graph.Features, + schemas: schema_lib.FeatureSetSchema, + group_specs: List[temporal_util.TimeseriesGroupSpec], + feature_padding: Dict[str, padding_lib.FeaturePadding], + padding_value: Any = 0, +) -> in_memory_graph.Features: + """Pads/caps timeseries features and generates matching mask features for a single entity set.""" + _validate_group_sequence_lengths(group_specs, feature_padding, schemas) + new_values: in_memory_graph.Features = {} + + padded_features = { + feat: fp.max_timeseries_len + for feat, fp in feature_padding.items() + if fp.max_timeseries_len is not None + } + + for feature_name, val in values.items(): + # Skip features that do not have a timeseries sequence length set. + if feature_name not in padded_features: + new_values[feature_name] = val + continue + + feature_schema = schemas[feature_name] + seq_len = padded_features[feature_name] + dtype = feature_format.FEATURE_FORMAT_TO_NP_DTYPE[feature_schema.format] + feat_shape = temporal_util.get_timeseries_step_shape(feature_schema) + + feature_padding_val = padding_value + if ( + padding_value == 0 + and feature_schema.format == schema_lib.FeatureFormat.BYTES + ): + feature_padding_val = b"" + + padded_matrix, mask_matrix = _pad_and_cap_single_feature( + raw_series=val, + seq_len=seq_len, + feat_shape=feat_shape, + padding_value=feature_padding_val, + dtype=dtype, + is_static_shape=feature_schema.is_static_shape(), + ) + + new_values[feature_name] = padded_matrix + if feature_schema.semantic == schema_lib.FeatureSemantic.MASK: + continue + + ts_group = feature_schema.group or feature_name + mask_name = temporal_util.get_mask_feature_name(feature_name, schemas) + if mask_name is None: + mask_name = f"{ts_group}_mask" + + # Copy over mask matrix if it doesn't exist yet. Only store once per + # group. + if mask_name not in new_values: + new_values[mask_name] = mask_matrix + + return new_values + + +def pad_timeseries_graph( + graph: in_memory_graph.InMemoryGraph, + schema: schema_lib.GraphSchema, + padding: padding_lib.Padding, + padding_value: Any = 0, + schema_cache: Optional[temporal_util.TimeseriesSchemaCache] = None, +) -> in_memory_graph.InMemoryGraph: + """Pads and caps all timeseries features in a single graph.""" + if not temporal_util.schema_has_timeseries_features(schema): + return graph + + if schema_cache is None: + schema_cache = temporal_util.extract_timeseries_schema_cache(schema) + + new_node_sets = {} + for ns_name, ns_schema in schema.node_sets.items(): + ns_val = graph.node_sets[ns_name] + ts_specs = schema_cache.node_sets[ns_name] + if not ts_specs: + new_node_sets[ns_name] = ns_val + continue + new_vals = pad_timeseries_feature_set( + values=ns_val.features, + schemas=ns_schema.features, + group_specs=ts_specs, + feature_padding=padding.node_sets[ns_name].features, + padding_value=padding_value, + ) + new_node_sets[ns_name] = in_memory_graph.InMemoryNodeSet( + num_nodes=ns_val.num_nodes, features=new_vals + ) + + new_edge_sets = {} + for es_name, es_schema in schema.edge_sets.items(): + es_val = graph.edge_sets[es_name] + ts_specs = schema_cache.edge_sets[es_name] + if not ts_specs: + new_edge_sets[es_name] = es_val + continue + new_vals = pad_timeseries_feature_set( + values=es_val.features, + schemas=es_schema.features, + group_specs=ts_specs, + feature_padding=padding.edge_sets[es_name].features, + padding_value=padding_value, + ) + new_edge_sets[es_name] = in_memory_graph.InMemoryEdgeSet( + adjacency=es_val.adjacency, features=new_vals + ) + + return in_memory_graph.InMemoryGraph( + node_sets=new_node_sets, edge_sets=new_edge_sets + ) diff --git a/dgf/src/transform/timeseries_padding_test.py b/dgf/src/transform/timeseries_padding_test.py new file mode 100644 index 0000000..0bf1844 --- /dev/null +++ b/dgf/src/transform/timeseries_padding_test.py @@ -0,0 +1,453 @@ +# Copyright 2022 Google LLC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Tests for padding and capping timeseries sequence features.""" + +from absl.testing import absltest +from absl.testing import parameterized +from dgf.src.data import in_memory_graph +from dgf.src.data import padding as padding_lib +from dgf.src.data import schema as schema_lib +from dgf.src.transform import timeseries_padding +from dgf.src.util import test_util +import numpy as np + + +def _make_graph_and_schema( + values: dict[str, np.ndarray], + schemas: dict[str, schema_lib.FeatureSchema], + node_set_name: str = "hardware", + num_nodes: int = 1, + edge_values: dict[str, np.ndarray] | None = None, + edge_schemas: dict[str, schema_lib.FeatureSchema] | None = None, + edge_set_name: str = "edges", +) -> tuple[in_memory_graph.InMemoryGraph, schema_lib.GraphSchema]: + edge_sets = {} + edge_set_schemas = {} + if edge_values is not None and edge_schemas is not None: + edge_sets[edge_set_name] = in_memory_graph.InMemoryEdgeSet( + adjacency=np.array([[0], [0]]), + features=dict(edge_values), + ) + edge_set_schemas[edge_set_name] = schema_lib.EdgeSchema( + source=node_set_name, + target=node_set_name, + features=edge_schemas, + ) + + return ( + in_memory_graph.InMemoryGraph( + node_sets={ + node_set_name: in_memory_graph.InMemoryNodeSet( + num_nodes=num_nodes, features=dict(values) + ) + }, + edge_sets=edge_sets, + ), + schema_lib.GraphSchema( + node_sets={node_set_name: schema_lib.NodeSchema(features=schemas)}, + edge_sets=edge_set_schemas, + ), + ) + + +def _ts_schema( + fmt: schema_lib.FeatureFormat = schema_lib.FeatureFormat.FLOAT_32, + sem: schema_lib.FeatureSemantic = schema_lib.FeatureSemantic.NUMERICAL, + group: str | None = None, + is_creation_time: bool = False, + shape: schema_lib.Shape = (None,), +) -> schema_lib.FeatureSchema: + return schema_lib.FeatureSchema( + format=fmt, + semantic=sem, + is_timeseries=True, + group=group, + is_creation_time=is_creation_time, + shape=shape, + ) + + +class TimeseriesPaddingTest(parameterized.TestCase): + + def test_capping_and_padding_with_padding_config(self): + graph, schema = _make_graph_and_schema( + values={ + "time": np.array( + [np.array([10, 20, 30, 40, 50]), np.array([5, 15])], + dtype=np.object_, + ), + "signal": np.array( + [ + np.array([1.0, 2.0, 3.0, 4.0, 5.0], dtype=np.float32), + np.array([0.5, 1.5], dtype=np.float32), + ], + dtype=np.object_, + ), + "id": np.array([101, 102]), + }, + schemas={ + "time": _ts_schema( + fmt=schema_lib.FeatureFormat.INTEGER_64, + sem=schema_lib.FeatureSemantic.TIMESTAMP, + is_creation_time=True, + group="time", + ), + "signal": _ts_schema(group="time"), + "id": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.INTEGER_64, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + ), + }, + num_nodes=2, + ) + padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=2, + features={ + "time": padding_lib.FeaturePadding(max_timeseries_len=3), + "signal": padding_lib.FeaturePadding(max_timeseries_len=3), + }, + ) + }, + edge_sets={}, + ) + new_graph = timeseries_padding.pad_timeseries_graph( + graph=graph, schema=schema, padding=padding + ) + new_schema = timeseries_padding.pad_timeseries_schema( + schema=schema, padding=padding + ) + hw_val = new_graph.node_sets["hardware"] + hw_sch = new_schema.node_sets["hardware"] + + expected_features = { + "time": np.array([[30, 40, 50], [0, 5, 15]], dtype=np.int64), + "signal": np.array( + [[3.0, 4.0, 5.0], [0.0, 0.5, 1.5]], dtype=np.float32 + ), + "time_mask": np.array([[True, True, True], [False, True, True]]), + "id": np.array([101, 102]), + } + test_util.assert_are_equal(self, hw_val.features, expected_features) + + self.assertEqual(hw_sch.features["time"].shape, (3,)) + self.assertTrue(hw_sch.features["time"].is_timeseries) + self.assertTrue(hw_sch.features["time_mask"].is_timeseries) + self.assertEqual( + hw_sch.features["time_mask"].semantic, schema_lib.FeatureSemantic.MASK + ) + + def test_edge_sets_and_non_timeseries(self): + graph = in_memory_graph.InMemoryGraph( + node_sets={ + "user": in_memory_graph.InMemoryNodeSet( + num_nodes=1, features={"age": np.array([30], dtype=np.int64)} + ), + }, + edge_sets={ + "clicks": in_memory_graph.InMemoryEdgeSet( + adjacency=np.array([[0], [0]], dtype=np.int64), + features={ + "time": np.array([np.array([100, 200])], dtype=np.object_) + }, + ), + "static_edge": in_memory_graph.InMemoryEdgeSet( + adjacency=np.array([[0], [0]], dtype=np.int64), + features={"weight": np.array([1.0], dtype=np.float32)}, + ), + }, + ) + schema = schema_lib.GraphSchema( + node_sets={ + "user": schema_lib.NodeSchema( + features={ + "age": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.INTEGER_64, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + ) + } + ) + }, + edge_sets={ + "clicks": schema_lib.EdgeSchema( + source="user", + target="user", + features={ + "time": _ts_schema( + fmt=schema_lib.FeatureFormat.INTEGER_64, + sem=schema_lib.FeatureSemantic.TIMESTAMP, + group="time", + ) + }, + ), + "static_edge": schema_lib.EdgeSchema( + source="user", + target="user", + features={ + "weight": schema_lib.FeatureSchema( + format=schema_lib.FeatureFormat.FLOAT_32, + semantic=schema_lib.FeatureSemantic.NUMERICAL, + ) + }, + ), + }, + ) + padding = padding_lib.Padding( + node_sets={}, + edge_sets={ + "clicks": padding_lib.EdgeSetPadding( + num_edges=1, + features={ + "time": padding_lib.FeaturePadding(max_timeseries_len=2) + }, + ) + }, + ) + new_graph = timeseries_padding.pad_timeseries_graph( + graph=graph, schema=schema, padding=padding + ) + new_schema = timeseries_padding.pad_timeseries_schema( + schema=schema, padding=padding + ) + + np.testing.assert_array_equal( + new_graph.node_sets["user"].features["age"], [30] + ) + np.testing.assert_array_equal( + new_graph.edge_sets["clicks"].features["time"][0], [100, 200] + ) + self.assertTrue( + new_schema.edge_sets["clicks"].features["time"].is_timeseries + ) + + def test_multidimensional_sequence(self): + graph, schema = _make_graph_and_schema( + values={ + "emb": np.array( + [np.array([[1.0, 2.0], [3.0, 4.0]], dtype=np.float32)], + dtype=np.object_, + ) + }, + schemas={ + "emb": _ts_schema( + sem=schema_lib.FeatureSemantic.EMBEDDING, + group="emb", + shape=(None, 2), + ) + }, + ) + padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=1, + features={ + "emb": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + new_graph = timeseries_padding.pad_timeseries_graph( + graph=graph, schema=schema, padding=padding + ) + new_schema = timeseries_padding.pad_timeseries_schema( + schema=schema, padding=padding + ) + hw_val = new_graph.node_sets["hardware"] + hw_sch = new_schema.node_sets["hardware"] + + expected_features = { + "emb": np.array( + [[[0.0, 0.0], [1.0, 2.0], [3.0, 4.0]]], dtype=np.float32 + ), + "emb_mask": np.array([[False, True, True]]), + } + test_util.assert_are_equal(self, hw_val.features, expected_features) + self.assertEqual(hw_sch.features["emb"].shape, (3, 2)) + self.assertEqual(hw_sch.features["emb_mask"].shape, (3,)) + self.assertEqual( + hw_sch.features["emb_mask"].semantic, schema_lib.FeatureSemantic.MASK + ) + + def test_empty_padding_returns_original_schema(self): + _, schema = _make_graph_and_schema( + values={ + "signal": np.array( + [np.array([1.0, 2.0], dtype=np.float32)], dtype=np.object_ + ), + }, + schemas={ + "signal": _ts_schema(group="sig", shape=(None,)), + }, + ) + empty_features_padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=1, + features={}, + ) + }, + edge_sets={}, + ) + # Empty feature padding returns original schema unchanged + schema_empty_features = timeseries_padding.pad_timeseries_schema( + schema, padding=empty_features_padding + ) + self.assertEqual(schema_empty_features, schema) + + def test_inconsistent_group_padding_raises_error(self): + _, schema = _make_graph_and_schema( + values={ + "time": np.array([np.array([1, 2])], dtype=np.object_), + "signal": np.array([np.array([1.0, 2.0])], dtype=np.object_), + }, + schemas={ + "time": _ts_schema( + fmt=schema_lib.FeatureFormat.INTEGER_64, + sem=schema_lib.FeatureSemantic.TIMESTAMP, + group="sensor", + ), + "signal": _ts_schema( + fmt=schema_lib.FeatureFormat.FLOAT_32, + group="sensor", + ), + }, + ) + inconsistent_padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + features={ + "time": padding_lib.FeaturePadding(max_timeseries_len=5), + }, + ) + }, + edge_sets={}, + ) + with self.assertRaisesRegex( + ValueError, "inconsistent padding configuration" + ): + timeseries_padding.pad_timeseries_schema( + schema, padding=inconsistent_padding + ) + + def test_conflicting_group_sequence_lengths_raises_error(self): + _, schema = _make_graph_and_schema( + values={ + "time": np.array([np.array([1, 2])], dtype=np.object_), + "signal": np.array([np.array([1.0, 2.0])], dtype=np.object_), + }, + schemas={ + "time": _ts_schema( + fmt=schema_lib.FeatureFormat.INTEGER_64, + sem=schema_lib.FeatureSemantic.TIMESTAMP, + group="sensor", + ), + "signal": _ts_schema( + fmt=schema_lib.FeatureFormat.FLOAT_32, + group="sensor", + ), + }, + ) + conflicting_padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=1, + features={ + "time": padding_lib.FeaturePadding(max_timeseries_len=5), + "signal": padding_lib.FeaturePadding(max_timeseries_len=10), + }, + ) + }, + edge_sets={}, + ) + with self.assertRaisesRegex(ValueError, "conflicting max_timeseries_len"): + timeseries_padding.pad_timeseries_schema( + schema, padding=conflicting_padding + ) + + + def test_bytes_timeseries_padding(self): + graph, schema = _make_graph_and_schema( + values={ + "tag": np.array( + [np.array([b"alpha", b"beta"]), np.array([b"gamma"])], + dtype=np.object_, + ), + }, + schemas={ + "tag": _ts_schema( + fmt=schema_lib.FeatureFormat.BYTES, + sem=schema_lib.FeatureSemantic.CATEGORICAL, + group="tag", + ), + }, + num_nodes=2, + ) + padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=2, + features={ + "tag": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + new_graph = timeseries_padding.pad_timeseries_graph( + graph=graph, schema=schema, padding=padding + ) + hw_val = new_graph.node_sets["hardware"] + np.testing.assert_array_equal( + hw_val.features["tag"], + np.array([[b"", b"alpha", b"beta"], [b"", b"", b"gamma"]]), + ) + np.testing.assert_array_equal( + hw_val.features["tag_mask"], + np.array([[False, True, True], [False, False, True]]), + ) + + def test_empty_entities_timeseries_padding(self): + graph, schema = _make_graph_and_schema( + values={ + "sig": np.empty((0,), dtype=np.object_), + }, + schemas={ + "sig": _ts_schema(group="sig"), + }, + num_nodes=0, + ) + padding = padding_lib.Padding( + node_sets={ + "hardware": padding_lib.NodeSetPadding( + num_nodes=0, + features={ + "sig": padding_lib.FeaturePadding(max_timeseries_len=3) + }, + ) + }, + edge_sets={}, + ) + new_graph = timeseries_padding.pad_timeseries_graph( + graph=graph, schema=schema, padding=padding + ) + hw_val = new_graph.node_sets["hardware"] + self.assertEqual(hw_val.features["sig"].shape, (0, 3)) + self.assertEqual(hw_val.features["sig_mask"].shape, (0, 3)) + + +if __name__ == "__main__": + absltest.main()