Class LogStreamManager
- java.lang.Object
-
- org.nuxeo.lib.stream.computation.log.LogStreamManager
-
- All Implemented Interfaces:
StreamManager
public class LogStreamManager extends Object implements StreamManager
StreamManager based on a LogManager- Since:
- 11.1
-
-
Field Summary
Fields Modifier and Type Field Description protected Map<Name,RecordFilterChain>
filters
static Codec<Record>
INTERNAL_CODEC
protected LogManager
logManager
static String
METRICS_STREAM
static String
PROCESSORS_STREAM
protected Map<String,Settings>
settings
protected Set<Name>
streams
protected Map<String,Topology>
topologies
-
Constructor Summary
Constructors Constructor Description LogStreamManager(LogManager logManager)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description LogOffset
append(String streamUrn, Record record)
Appends a record to a processor's source stream.StreamProcessor
createStreamProcessor(String processorName)
Creates a registered processor without starting it.LogTailer<Record>
createTailer(Name computationName, Collection<LogPartition> streamPartitions)
protected Codec<Record>
getCodec(Collection<Name> streams)
RecordFilter
getFilter(Name stream)
LogManager
getLogManager()
protected Map<String,String>
getSystemMetadata()
protected void
initAppenders(Collection<String> streams, Settings settings)
protected void
initInternalStream(Name stream)
protected void
initInternalStreams()
protected void
initStream(String streamName, Settings settings)
protected void
initStreams(Topology topology, Settings settings)
void
register(String processorName, Topology topology, Settings settings)
Registers a processor and initializes the underlying streams, this is needed before creating a processor or appending record in source streams.void
register(List<String> streams, Settings settings)
Registers some source Streams without any processors.protected void
registerFilters(Collection<String> streams, Settings settings)
LogTailer<Record>
subscribe(Name computationName, Collection<Name> streams, RebalanceListener listener)
boolean
supportSubscribe()
-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
-
Methods inherited from interface org.nuxeo.lib.stream.computation.StreamManager
registerAndCreateProcessor
-
-
-
-
Field Detail
-
PROCESSORS_STREAM
public static final String PROCESSORS_STREAM
- See Also:
- Constant Field Values
-
METRICS_STREAM
public static final String METRICS_STREAM
- See Also:
- Constant Field Values
-
INTERNAL_CODEC
public static final Codec<Record> INTERNAL_CODEC
-
logManager
protected final LogManager logManager
-
topologies
protected final Map<String,Topology> topologies
-
filters
protected final Map<Name,RecordFilterChain> filters
-
-
Constructor Detail
-
LogStreamManager
public LogStreamManager(LogManager logManager)
-
-
Method Detail
-
initInternalStreams
protected void initInternalStreams()
-
initInternalStream
protected void initInternalStream(Name stream)
-
register
public void register(String processorName, Topology topology, Settings settings)
Description copied from interface:StreamManager
Registers a processor and initializes the underlying streams, this is needed before creating a processor or appending record in source streams.- Specified by:
register
in interfaceStreamManager
-
register
public void register(List<String> streams, Settings settings)
Description copied from interface:StreamManager
Registers some source Streams without any processors.- Specified by:
register
in interfaceStreamManager
-
createStreamProcessor
public StreamProcessor createStreamProcessor(String processorName)
Description copied from interface:StreamManager
Creates a registered processor without starting it.- Specified by:
createStreamProcessor
in interfaceStreamManager
-
getSystemMetadata
protected Map<String,String> getSystemMetadata()
-
getLogManager
public LogManager getLogManager()
-
append
public LogOffset append(String streamUrn, Record record)
Description copied from interface:StreamManager
Appends a record to a processor's source stream.- Specified by:
append
in interfaceStreamManager
-
supportSubscribe
public boolean supportSubscribe()
-
subscribe
public LogTailer<Record> subscribe(Name computationName, Collection<Name> streams, RebalanceListener listener)
-
createTailer
public LogTailer<Record> createTailer(Name computationName, Collection<LogPartition> streamPartitions)
-
getFilter
public RecordFilter getFilter(Name stream)
-
getCodec
protected Codec<Record> getCodec(Collection<Name> streams)
-
initStreams
protected void initStreams(Topology topology, Settings settings)
-
initStream
protected void initStream(String streamName, Settings settings)
-
initAppenders
protected void initAppenders(Collection<String> streams, Settings settings)
-
registerFilters
protected void registerFilters(Collection<String> streams, Settings settings)
-
-