Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@
import org.apache.samza.system.SystemFactory;
import org.apache.samza.system.SystemProducer;
import org.apache.samza.system.SystemStream;

import org.apache.samza.system.kinesis.consumer.KinesisSystemConsumer;


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@
import org.apache.samza.system.eventhub.producer.EventHubSystemProducer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.collection.JavaConversions;

import java.time.Duration;
import java.util.HashMap;
Expand Down Expand Up @@ -109,7 +108,7 @@ public EventHubConfig(Config config) {
StreamConfig streamConfig = new StreamConfig(config);

LOG.info("Building mappings from physicalName to streamId");
JavaConversions.asJavaCollection(streamConfig.getStreamIds())
streamConfig.getStreamIds()
.forEach((streamId) -> {
String physicalName = streamConfig.getPhysicalName(streamId);
LOG.info("Obtained physicalName: {} for streamId: {} ", physicalName, streamId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -160,15 +160,15 @@ public String getSystemOffsetDefault(String systemName) {
* @return the key serde for the {@code systemName}, or empty if it was not found
*/
public Optional<String> getSystemKeySerde(String systemName) {
return getSystemDefaultStreamProperty(systemName, StreamConfig.KEY_SERDE());
return getSystemDefaultStreamProperty(systemName, StreamConfig.KEY_SERDE);
}

/**
* @param systemName name of the system
* @return the message serde for the {@code systemName}, or empty if it was not found
*/
public Optional<String> getSystemMsgSerde(String systemName) {
return getSystemDefaultStreamProperty(systemName, StreamConfig.MSG_SERDE());
return getSystemDefaultStreamProperty(systemName, StreamConfig.MSG_SERDE);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,7 @@ private void configureTables(Map<String, String> generatedConfig, Config origina
.map(sideInput -> StreamUtil.getSystemStreamFromNameOrId(originalConfig, sideInput))
.forEach(systemStream -> {
inputs.add(StreamUtil.getNameFromSystemStream(systemStream));
generatedConfig.put(String.format(StreamConfig.STREAM_PREFIX() + StreamConfig.BOOTSTRAP(),
generatedConfig.put(String.format(StreamConfig.STREAM_PREFIX + StreamConfig.BOOTSTRAP,
systemStream.getSystem(), systemStream.getStream()), "true");
});
}
Expand Down Expand Up @@ -314,14 +314,14 @@ private void configureSerdes(Map<String, String> configs, Map<String, StreamEdge

// set key and msg serdes for streams to the serde names generated above
streamKeySerdes.forEach((streamId, serde) -> {
String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX(), streamId);
String keySerdeConfigKey = streamIdPrefix + StreamConfig.KEY_SERDE();
String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX, streamId);
String keySerdeConfigKey = streamIdPrefix + StreamConfig.KEY_SERDE;
configs.put(keySerdeConfigKey, serdeUUIDs.get(serde));
});

streamMsgSerdes.forEach((streamId, serde) -> {
String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX(), streamId);
String valueSerdeConfigKey = streamIdPrefix + StreamConfig.MSG_SERDE();
String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX, streamId);
String valueSerdeConfigKey = streamIdPrefix + StreamConfig.MSG_SERDE;
configs.put(valueSerdeConfigKey, serdeUUIDs.get(serde));
});

Expand Down
24 changes: 12 additions & 12 deletions samza-core/src/main/java/org/apache/samza/execution/StreamEdge.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,6 @@

package org.apache.samza.execution;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

import org.apache.samza.config.ApplicationConfig;
import org.apache.samza.config.Config;
import org.apache.samza.config.MapConfig;
Expand All @@ -32,6 +27,11 @@
import org.apache.samza.system.SystemStream;
import org.apache.samza.util.StreamUtil;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;


/**
* A StreamEdge connects the source {@link JobNode}s to the target {@link JobNode}s with a stream.
Expand Down Expand Up @@ -122,21 +122,21 @@ Config generateConfig() {
Map<String, String> streamConfig = new HashMap<>();
StreamSpec spec = getStreamSpec();
String streamId = spec.getId();
streamConfig.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), streamId), spec.getSystemName());
streamConfig.put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), streamId), spec.getPhysicalName());
streamConfig.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, streamId), spec.getSystemName());
streamConfig.put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, streamId), spec.getPhysicalName());
if (isIntermediate()) {
streamConfig.put(String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID(), streamId), "true");
streamConfig.put(String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID(), streamId), "true");
streamConfig.put(String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID, streamId), "true");
streamConfig.put(String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, streamId), "true");

// Setting offset.default to oldest only if the job is running in batch mode
if (ApplicationConfig.ApplicationMode.BATCH.equals(new ApplicationConfig(config).getAppMode())) {
streamConfig.put(String.format(StreamConfig.CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID(), streamId), "oldest");
streamConfig.put(String.format(StreamConfig.CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID, streamId), "oldest");
}

streamConfig.put(String.format(StreamConfig.PRIORITY_FOR_STREAM_ID(), streamId), String.valueOf(Integer.MAX_VALUE));
streamConfig.put(String.format(StreamConfig.PRIORITY_FOR_STREAM_ID, streamId), String.valueOf(Integer.MAX_VALUE));
}
spec.getConfig().forEach((property, value) -> {
streamConfig.put(String.format(StreamConfig.STREAM_ID_PREFIX(), streamId) + property, value);
streamConfig.put(String.format(StreamConfig.STREAM_ID_PREFIX, streamId) + property, value);
});

return new MapConfig(streamConfig);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,6 @@
import org.apache.samza.util.StreamUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.collection.JavaConversions;


public class StreamManager {
Expand Down Expand Up @@ -113,7 +112,7 @@ public void clearStreamsFromPreviousRun(Config prevConfig, ClassLoader classLoad
StreamConfig streamConfig = new StreamConfig(prevConfig);

//Find all intermediate streams and clean up
Set<StreamSpec> intStreams = JavaConversions.asJavaCollection(streamConfig.getStreamIds()).stream()
Set<StreamSpec> intStreams = streamConfig.getStreamIds().stream()
.filter(streamConfig::getIsIntermediateStream)
.map(id -> new StreamSpec(id, streamConfig.getPhysicalName(id), streamConfig.getSystem(id)))
.collect(Collectors.toSet());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,7 @@ static Multimap<SystemStream, String> getStreamToConsumerTasks(JobModel jobModel
* @return mapping from output streams to input streams
*/
static Multimap<SystemStream, SystemStream> getIntermediateToInputStreamsMap(
OperatorSpecGraph specGraph, StreamConfig streamConfig) {
OperatorSpecGraph specGraph, StreamConfig streamConfig) {
Multimap<SystemStream, SystemStream> outputToInputStreams = HashMultimap.create();
specGraph.getInputOperators().entrySet().stream()
.forEach(entry -> {
Expand All @@ -391,7 +391,7 @@ static Multimap<SystemStream, SystemStream> getIntermediateToInputStreamsMap(
}

private static void computeOutputToInput(SystemStream input, OperatorSpec opSpec,
Multimap<SystemStream, SystemStream> outputToInputStreams, StreamConfig streamConfig) {
Multimap<SystemStream, SystemStream> outputToInputStreams, StreamConfig streamConfig) {
if (opSpec instanceof PartitionByOperatorSpec) {
PartitionByOperatorSpec spec = (PartitionByOperatorSpec) opSpec;
SystemStream systemStream = streamConfig.streamIdToSystemStream(spec.getOutputStream().getStreamId());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,6 @@
package org.apache.samza.storage;

import com.google.common.collect.ImmutableMap;
import java.io.File;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.Set;

import java.util.stream.Collectors;
import org.apache.samza.clustermanager.StandbyTaskUtil;
import org.apache.samza.container.TaskName;
import org.apache.samza.job.model.TaskMode;
Expand All @@ -42,6 +35,13 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.File;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;


public class StorageManagerUtil {
private static final Logger LOG = LoggerFactory.getLogger(StorageManagerUtil.class);
Expand All @@ -57,8 +57,8 @@ public class StorageManagerUtil {
/**
* Fetch the starting offset for the input {@link SystemStreamPartition}
*
* Note: The method doesn't respect {@link org.apache.samza.config.StreamConfig#CONSUMER_OFFSET_DEFAULT()} and
* {@link org.apache.samza.config.StreamConfig#CONSUMER_RESET_OFFSET()} configurations. It will use the locally
* Note: The method doesn't respect {@link org.apache.samza.config.StreamConfig#CONSUMER_OFFSET_DEFAULT} and
* {@link org.apache.samza.config.StreamConfig#CONSUMER_RESET_OFFSET} configurations. It will use the locally
* checkpointed offset if it is valid, or fall back to oldest offset of the stream.
*
* @param ssp system stream partition for which starting offset is requested
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,18 +20,6 @@
package org.apache.samza.storage;

import com.google.common.annotations.VisibleForTesting;
import java.io.File;
import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.apache.samza.Partition;
import org.apache.samza.SamzaException;
import org.apache.samza.config.Config;
Expand All @@ -51,6 +39,19 @@
import org.slf4j.LoggerFactory;
import scala.collection.JavaConverters;

import java.io.File;
import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.stream.Collectors;


/**
* A storage manager for all side input stores. It is associated with each {@link org.apache.samza.container.TaskInstance}
Expand Down Expand Up @@ -162,8 +163,8 @@ public StorageEngine getStore(String storeName) {
/**
* Gets the starting offset for the given side input {@link SystemStreamPartition}.
*
* Note: The method doesn't respect {@link org.apache.samza.config.StreamConfig#CONSUMER_OFFSET_DEFAULT()} and
* {@link org.apache.samza.config.StreamConfig#CONSUMER_RESET_OFFSET()} configurations. It will use the local offset
* Note: The method doesn't respect {@link org.apache.samza.config.StreamConfig#CONSUMER_OFFSET_DEFAULT} and
* {@link org.apache.samza.config.StreamConfig#CONSUMER_RESET_OFFSET} configurations. It will use the local offset
* file if it is valid, else it will fall back to oldest offset in the stream.
*
* @param ssp side input system stream partition to get the starting offset for
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,16 @@
*/
package org.apache.samza.util;

import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;
import org.apache.samza.SamzaException;
import org.apache.samza.config.Config;
import org.apache.samza.config.StreamConfig;
import org.apache.samza.system.StreamSpec;
import org.apache.samza.system.SystemStream;

import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;

public class StreamUtil {
/**
* Gets the {@link SystemStream} corresponding to the provided stream, which may be
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,7 @@ import java.util.concurrent.ConcurrentHashMap
import org.apache.commons.lang3.StringUtils
import org.apache.samza.SamzaException
import org.apache.samza.annotation.InterfaceStability
import org.apache.samza.config.StreamConfig.Config2Stream
import org.apache.samza.config.{Config, SystemConfig}
import org.apache.samza.config.{Config, StreamConfig, SystemConfig}
import org.apache.samza.container.TaskName
import org.apache.samza.startpoint.{Startpoint, StartpointManager}
import org.apache.samza.system.SystemStreamMetadata.OffsetType
Expand Down Expand Up @@ -80,13 +79,15 @@ object OffsetManager extends Logging {
offsetManagerMetrics: OffsetManagerMetrics = new OffsetManagerMetrics) = {
debug("Building offset manager for %s." format systemStreamMetadata)

val streamConfig = new StreamConfig(config)

val offsetSettings = systemStreamMetadata
.map {
case (systemStream, systemStreamMetadata) =>
// Get default offset.
val streamDefaultOffset = config.getDefaultStreamOffset(systemStream)
val streamDefaultOffset = streamConfig.getDefaultStreamOffset(systemStream)
val systemDefaultOffset = new SystemConfig(config).getSystemOffsetDefault(systemStream.getSystem)
val defaultOffsetType = if (streamDefaultOffset.isDefined) {
val defaultOffsetType = if (streamDefaultOffset.isPresent) {
OffsetType.valueOf(streamDefaultOffset.get.toUpperCase)
} else if (systemDefaultOffset != null) {
OffsetType.valueOf(systemDefaultOffset.toUpperCase)
Expand All @@ -97,7 +98,7 @@ object OffsetManager extends Logging {
debug("Using default offset %s for %s." format (defaultOffsetType, systemStream))

// Get reset offset.
val resetOffset = config.getResetOffset(systemStream)
val resetOffset = streamConfig.getResetOffset(systemStream)
debug("Using reset offset %s for %s." format (resetOffset, systemStream))

// Build OffsetSetting so we can create a map for OffsetManager.
Expand Down
Loading