Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
211 changes: 163 additions & 48 deletions src/mlpro/oa/streams/tasks/clusteranalyzers/basics.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,35 +51,40 @@
## -- 2025-09-03 1.8.0 DA Class ClusterAnalyzer:
## -- - Bugfix: added missing parameter p_thrs_cluster_influence
## -- - Method _get_cluster_relations(): robustness for negative influences
## -- 2026-08-05 1.9.0 DA New classes ClusterActions, ClusterInfrastructure
## -------------------------------------------------------------------------------------------------

"""
Ver. 1.8.0 (2025-09-03)
Ver. 1.9.0 (2026-08-05)

This module provides a template class for online cluster analysis.
"""


from typing import List, Tuple
from abc import ABC, abstractmethod

try:
from matplotlib.figure import Figure
except:
class Figure : pass

from mlpro.bf.events import Event as MLProEvent
from mlpro.bf.events import Event as MLProEvent, EventManager
from mlpro.bf.math.properties import *
from mlpro.bf.mt import PlotSettings
from mlpro.bf.streams import Instance, InstDict
from mlpro.bf.various import *
from mlpro.bf.plot import *

from mlpro.oa.streams import OAStreamTask
from mlpro.bf.math.normalizers import Normalizer
from mlpro.oa.streams.tasks.clusteranalyzers.clusters import Cluster, ClusterId



# Export list for public API
__all__ = [ 'ClusterAnalyzer',
__all__ = [ 'ClusterActions',
'ClusterInfrastructure',
'ClusterAnalyzer',
'ClusterId',
'ResultItem' ]

Expand All @@ -91,7 +96,69 @@ class Figure : pass

## -------------------------------------------------------------------------------------------------
## -------------------------------------------------------------------------------------------------
class ClusterAnalyzer (OAStreamTask):
class ClusterActions (ABC):

# Possible result scopes for methods get_cluster_memberships() and get_cluster_influences()
C_RESULT_SCOPE_ALL : int = 0
C_RESULT_SCOPE_MAX : int = 1


## -------------------------------------------------------------------------------------------------
def __init__( self ):
self._clusters = {}


## -------------------------------------------------------------------------------------------------
@abstractmethod
def get_cluster_memberships( self,
p_instance : Instance,
p_scope : int = C_RESULT_SCOPE_MAX ) -> List[ResultItem]:
"""
Method to determine the relative membership of the given instance to each cluster as a value
in [0,1].

See also: method Cluster.get_membership().

Parameters
----------
p_instance : Instance
Instance to be evaluated.
p_scope : int
Scope of the result list. See class attributes C_RESULT_SCOPE_* for possible values. Default
value is C_RESULT_SCOPE_MAX.

Returns
-------
List[ResultItem]
List of membership items which are tuples of a cluster id, a relative membership value
in [0,1], and a reference to the cluster object.
"""

...


## -------------------------------------------------------------------------------------------------
@abstractmethod
def get_cluster_influences( self,
p_instance : Instance,
p_scope : int = C_RESULT_SCOPE_MAX ) -> List[ResultItem]: ...


## -------------------------------------------------------------------------------------------------
def _get_clusters(self):
return self._clusters


## -------------------------------------------------------------------------------------------------
clusters = property( fget = _get_clusters )





## -------------------------------------------------------------------------------------------------
## -------------------------------------------------------------------------------------------------
class ClusterInfrastructure (EventManager, Plottable, ClusterActions):
"""
Base class for online cluster analysis. It raises an event when a cluster was added or removed.

Expand All @@ -104,14 +171,6 @@ class ClusterAnalyzer (OAStreamTask):

Parameters
----------
p_name : str
Optional name of the task. Default is None.
p_range_max : int
Maximum range of asynchonicity. See class Range. Default is Range.C_RANGE_PROCESS.
p_ada : bool
Boolean switch for adaptivitiy. Default = True.
p_duplicate_data : bool
If True, instances will be duplicated before processing. Default = False.
p_cls_cluster
Cluster class (Class Cluster or a child class).
p_cluster_limit : int
Expand All @@ -122,8 +181,6 @@ class ClusterAnalyzer (OAStreamTask):
Boolean switch for visualisation. Default = False.
p_logging
Log level (see constants of class Log). Default: Log.C_LOG_ALL
p_kwargs : dict
Further optional named parameters.

Attributes
----------
Expand All @@ -136,19 +193,14 @@ class ClusterAnalyzer (OAStreamTask):
are handed over to each new cluster.
"""

C_TYPE = 'Cluster Analyzer'
C_TYPE = 'Cluster Compound'

C_EVENT_CLUSTER_ADDED = 'CLUSTER_ADDED'
C_EVENT_CLUSTER_REMOVED = 'CLUSTER_REMOVED'

C_PLOT_ACTIVE = True
C_PLOT_STANDALONE = False

# Possible result scopes for methods get_cluster_memberships() and get_cluster_influences()
C_RESULT_SCOPE_ALL : int = 0
# C_RESULT_SCOPE_NONZERO : int = 1
C_RESULT_SCOPE_MAX : int = 2

# List of cluster properties supported/maintained by the algorithm
C_CLUSTER_PROPERTIES : PropertyDefinitions = []

Expand All @@ -157,26 +209,16 @@ class ClusterAnalyzer (OAStreamTask):

## -------------------------------------------------------------------------------------------------
def __init__( self,
p_name: str = None,
p_range_max = OAStreamTask.C_RANGE_THREAD,
p_ada: bool = True,
p_duplicate_data: bool = False,
p_cls_cluster : type = Cluster,
p_cluster_limit : int = 0,
p_thrs_cluster_influence : float = None,
p_visualize: bool = False,
p_logging = Log.C_LOG_ALL,
**p_kwargs ):

super().__init__( p_name = p_name,
p_range_max = p_range_max,
p_ada = p_ada,
p_duplicate_data = p_duplicate_data,
p_visualize = p_visualize,
p_logging = p_logging,
**p_kwargs )

self._clusters = {}
p_logging = Log.C_LOG_WE ):

EventManager.__init__( self, p_logging = p_logging )
Plottable.__init__( self, p_visualize = p_visualize )
ClusterActions.__init__( self )

self._next_cluster_id : ClusterId = -1

self._cls_cluster = p_cls_cluster
Expand Down Expand Up @@ -219,11 +261,6 @@ def align_cluster_properties( self, p_properties : PropertyDefinitions ) -> list
return unknown_properties


## -------------------------------------------------------------------------------------------------
def _run(self, p_instances : InstDict):
self.adapt( p_instances = p_instances )


## -------------------------------------------------------------------------------------------------
def new_cluster_allowed(self) -> bool:
"""
Expand All @@ -243,11 +280,6 @@ def get_cluster_cls(self):
return self._cls_cluster


## -------------------------------------------------------------------------------------------------
def _get_clusters(self):
return self._clusters


## -------------------------------------------------------------------------------------------------
def _get_next_cluster_id(self) -> ClusterId:
self._next_cluster_id += 1
Expand Down Expand Up @@ -505,5 +537,88 @@ def _renormalize(self, p_normalizer: Normalizer):
cluster.renormalize( p_normalizer=p_normalizer )





## -------------------------------------------------------------------------------------------------
## -------------------------------------------------------------------------------------------------
clusters = property( fget = _get_clusters )
class ClusterAnalyzer (OAStreamTask, ClusterInfrastructure):
"""
Base class for online cluster analysis. It raises an event when a cluster was added or removed.

Steps to implement a new algorithm are:
- Create a new class and inherit from this base class
- Specify all cluster properties provided/maintained by your algorithm in C_CLUSTER_PROPERTIES.
- Implement method self._adapt() to update your cluster list on new instances
- Implement method self._adapt_reverse() to update your cluster list on obsolete instances
- New cluster: hand over self._cluster_properties.values() on instantiation

Parameters
----------
p_name : str
Optional name of the task. Default is None.
p_range_max : int
Maximum range of asynchonicity. See class Range. Default is Range.C_RANGE_PROCESS.
p_ada : bool
Boolean switch for adaptivitiy. Default = True.
p_duplicate_data : bool
If True, instances will be duplicated before processing. Default = False.
p_cls_cluster
Cluster class (Class Cluster or a child class).
p_cluster_limit : int
Optional limit for clusters to be created. Default = 0 (no limit).
p_thrs_cluster_influence : float
Threshold for cluster influence. Default = 0.0.
p_visualize : bool
Boolean switch for visualisation. Default = False.
p_logging
Log level (see constants of class Log). Default: Log.C_LOG_ALL
p_kwargs : dict
Further optional named parameters.

Attributes
----------
C_RESULT_SCOPE_ALL : int = 0
Result scope, that includes all clusters
C_RESULT_SCOPE_MAX : int = 2
Result scope, that includes just the cluster with the highest result value.
C_CLUSTER_PROPERTIES : PropertyDefinitions
List of cluster properties supported/maintained by the algorithm. These properties
are handed over to each new cluster.
"""

C_TYPE = 'Cluster Analyzer'

## -------------------------------------------------------------------------------------------------
def __init__( self,
p_name: str = None,
p_range_max = OAStreamTask.C_RANGE_THREAD,
p_ada: bool = True,
p_duplicate_data: bool = False,
p_cls_cluster : type = Cluster,
p_cluster_limit : int = 0,
p_thrs_cluster_influence : float = None,
p_visualize: bool = False,
p_logging = Log.C_LOG_ALL,
**p_kwargs ):

OAStreamTask.__init__( self,
p_name = p_name,
p_range_max = p_range_max,
p_ada = p_ada,
p_duplicate_data = p_duplicate_data,
p_visualize = False,
p_logging = Log.C_LOG_NOTHING,
**p_kwargs )

ClusterInfrastructure.__init__( self,
p_cls_cluster = p_cls_cluster,
p_cluster_limit = p_cluster_limit,
p_thrs_cluster_influence = p_thrs_cluster_influence,
p_visualize = p_visualize,
p_logging = p_logging )


## -------------------------------------------------------------------------------------------------
def _run(self, p_instances : InstDict):
self.adapt( p_instances = p_instances )