diff --git a/samza-aws/src/main/java/org/apache/samza/system/kinesis/KinesisSystemFactory.java b/samza-aws/src/main/java/org/apache/samza/system/kinesis/KinesisSystemFactory.java index 0086643b6b..27580220c4 100644 --- a/samza-aws/src/main/java/org/apache/samza/system/kinesis/KinesisSystemFactory.java +++ b/samza-aws/src/main/java/org/apache/samza/system/kinesis/KinesisSystemFactory.java @@ -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; diff --git a/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java b/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java index d6cdab560f..69a921e1a3 100644 --- a/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java +++ b/samza-azure/src/main/java/org/apache/samza/system/eventhub/EventHubConfig.java @@ -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; @@ -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); diff --git a/samza-core/src/main/java/org/apache/samza/config/SystemConfig.java b/samza-core/src/main/java/org/apache/samza/config/SystemConfig.java index 1403b0db17..bfbc2979e9 100644 --- a/samza-core/src/main/java/org/apache/samza/config/SystemConfig.java +++ b/samza-core/src/main/java/org/apache/samza/config/SystemConfig.java @@ -160,7 +160,7 @@ public String getSystemOffsetDefault(String systemName) { * @return the key serde for the {@code systemName}, or empty if it was not found */ public Optional getSystemKeySerde(String systemName) { - return getSystemDefaultStreamProperty(systemName, StreamConfig.KEY_SERDE()); + return getSystemDefaultStreamProperty(systemName, StreamConfig.KEY_SERDE); } /** @@ -168,7 +168,7 @@ public Optional getSystemKeySerde(String systemName) { * @return the message serde for the {@code systemName}, or empty if it was not found */ public Optional getSystemMsgSerde(String systemName) { - return getSystemDefaultStreamProperty(systemName, StreamConfig.MSG_SERDE()); + return getSystemDefaultStreamProperty(systemName, StreamConfig.MSG_SERDE); } /** diff --git a/samza-core/src/main/java/org/apache/samza/execution/JobNodeConfigurationGenerator.java b/samza-core/src/main/java/org/apache/samza/execution/JobNodeConfigurationGenerator.java index 7e94667c57..ca46d1c3e1 100644 --- a/samza-core/src/main/java/org/apache/samza/execution/JobNodeConfigurationGenerator.java +++ b/samza-core/src/main/java/org/apache/samza/execution/JobNodeConfigurationGenerator.java @@ -251,7 +251,7 @@ private void configureTables(Map 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"); }); } @@ -314,14 +314,14 @@ private void configureSerdes(Map configs, Map { - 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)); }); diff --git a/samza-core/src/main/java/org/apache/samza/execution/StreamEdge.java b/samza-core/src/main/java/org/apache/samza/execution/StreamEdge.java index 80d1bad626..4999d06993 100644 --- a/samza-core/src/main/java/org/apache/samza/execution/StreamEdge.java +++ b/samza-core/src/main/java/org/apache/samza/execution/StreamEdge.java @@ -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; @@ -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. @@ -122,21 +122,21 @@ Config generateConfig() { Map 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); diff --git a/samza-core/src/main/java/org/apache/samza/execution/StreamManager.java b/samza-core/src/main/java/org/apache/samza/execution/StreamManager.java index fe0d9fd0aa..8d6cb01d71 100644 --- a/samza-core/src/main/java/org/apache/samza/execution/StreamManager.java +++ b/samza-core/src/main/java/org/apache/samza/execution/StreamManager.java @@ -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 { @@ -113,7 +112,7 @@ public void clearStreamsFromPreviousRun(Config prevConfig, ClassLoader classLoad StreamConfig streamConfig = new StreamConfig(prevConfig); //Find all intermediate streams and clean up - Set intStreams = JavaConversions.asJavaCollection(streamConfig.getStreamIds()).stream() + Set intStreams = streamConfig.getStreamIds().stream() .filter(streamConfig::getIsIntermediateStream) .map(id -> new StreamSpec(id, streamConfig.getPhysicalName(id), streamConfig.getSystem(id))) .collect(Collectors.toSet()); diff --git a/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImplGraph.java b/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImplGraph.java index 2b953211d8..4cf7201937 100644 --- a/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImplGraph.java +++ b/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImplGraph.java @@ -380,7 +380,7 @@ static Multimap getStreamToConsumerTasks(JobModel jobModel * @return mapping from output streams to input streams */ static Multimap getIntermediateToInputStreamsMap( - OperatorSpecGraph specGraph, StreamConfig streamConfig) { + OperatorSpecGraph specGraph, StreamConfig streamConfig) { Multimap outputToInputStreams = HashMultimap.create(); specGraph.getInputOperators().entrySet().stream() .forEach(entry -> { @@ -391,7 +391,7 @@ static Multimap getIntermediateToInputStreamsMap( } private static void computeOutputToInput(SystemStream input, OperatorSpec opSpec, - Multimap outputToInputStreams, StreamConfig streamConfig) { + Multimap outputToInputStreams, StreamConfig streamConfig) { if (opSpec instanceof PartitionByOperatorSpec) { PartitionByOperatorSpec spec = (PartitionByOperatorSpec) opSpec; SystemStream systemStream = streamConfig.streamIdToSystemStream(spec.getOutputStream().getStreamId()); diff --git a/samza-core/src/main/java/org/apache/samza/storage/StorageManagerUtil.java b/samza-core/src/main/java/org/apache/samza/storage/StorageManagerUtil.java index 293b9ab910..a3f971257c 100644 --- a/samza-core/src/main/java/org/apache/samza/storage/StorageManagerUtil.java +++ b/samza-core/src/main/java/org/apache/samza/storage/StorageManagerUtil.java @@ -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; @@ -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); @@ -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 diff --git a/samza-core/src/main/java/org/apache/samza/storage/TaskSideInputStorageManager.java b/samza-core/src/main/java/org/apache/samza/storage/TaskSideInputStorageManager.java index 3de4f5d4a8..1bdd6b7931 100644 --- a/samza-core/src/main/java/org/apache/samza/storage/TaskSideInputStorageManager.java +++ b/samza-core/src/main/java/org/apache/samza/storage/TaskSideInputStorageManager.java @@ -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; @@ -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} @@ -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 diff --git a/samza-core/src/main/java/org/apache/samza/util/StreamUtil.java b/samza-core/src/main/java/org/apache/samza/util/StreamUtil.java index 96cdd6c4af..6f00549c32 100644 --- a/samza-core/src/main/java/org/apache/samza/util/StreamUtil.java +++ b/samza-core/src/main/java/org/apache/samza/util/StreamUtil.java @@ -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 diff --git a/samza-core/src/main/scala/org/apache/samza/checkpoint/OffsetManager.scala b/samza-core/src/main/scala/org/apache/samza/checkpoint/OffsetManager.scala index cdeddb05fb..238280a41a 100644 --- a/samza-core/src/main/scala/org/apache/samza/checkpoint/OffsetManager.scala +++ b/samza-core/src/main/scala/org/apache/samza/checkpoint/OffsetManager.scala @@ -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 @@ -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) @@ -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. diff --git a/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.java b/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.java new file mode 100644 index 0000000000..8ee044e85c --- /dev/null +++ b/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.java @@ -0,0 +1,396 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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 + * + * http://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. + */ + +package org.apache.samza.config; + +import com.google.common.collect.Sets; +import org.apache.samza.system.SystemStream; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.lang.invoke.MethodHandles; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +public class StreamConfig extends MapConfig { + public static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); + + + // Samza configs for streams + public static final String SAMZA_PROPERTY = "samza."; + public static final String SYSTEM = SAMZA_PROPERTY + "system"; + public static final String PHYSICAL_NAME = SAMZA_PROPERTY + "physical.name"; + public static final String MSG_SERDE = SAMZA_PROPERTY + "msg.serde"; + public static final String KEY_SERDE = SAMZA_PROPERTY + "key.serde"; + public static final String CONSUMER_RESET_OFFSET = SAMZA_PROPERTY + "reset.offset"; + public static final String CONSUMER_OFFSET_DEFAULT = SAMZA_PROPERTY + "offset.default"; + public static final String BOOTSTRAP = SAMZA_PROPERTY + "bootstrap"; + public static final String PRIORITY = SAMZA_PROPERTY + "priority"; + public static final String IS_INTERMEDIATE = SAMZA_PROPERTY + "intermediate"; + public static final String DELETE_COMMITTED_MESSAGES = SAMZA_PROPERTY + "delete.committed.messages"; + public static final String IS_BOUNDED = SAMZA_PROPERTY + "bounded"; + public static final String BROADCAST = SAMZA_PROPERTY + "broadcast"; + + // We don't want any external dependencies on these patterns while both exist. + // Use the corresponding get*() method to ensure proper values. + private static final String STREAMS_PREFIX = "streams."; + + public static final String STREAM_PREFIX = "systems.%s.streams.%s."; + public static final String STREAM_ID_PREFIX = STREAMS_PREFIX + "%s."; + public static final String SYSTEM_FOR_STREAM_ID = STREAM_ID_PREFIX + SYSTEM; + public static final String PHYSICAL_NAME_FOR_STREAM_ID = STREAM_ID_PREFIX + PHYSICAL_NAME; + public static final String IS_INTERMEDIATE_FOR_STREAM_ID = STREAM_ID_PREFIX + IS_INTERMEDIATE; + public static final String DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID = STREAM_ID_PREFIX + DELETE_COMMITTED_MESSAGES; + public static final String IS_BOUNDED_FOR_STREAM_ID = STREAM_ID_PREFIX + IS_BOUNDED; + public static final String PRIORITY_FOR_STREAM_ID = STREAM_ID_PREFIX + PRIORITY; + public static final String CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID = STREAM_ID_PREFIX + CONSUMER_OFFSET_DEFAULT; + public static final String BOOTSTRAP_FOR_STREAM_ID = STREAM_ID_PREFIX + BOOTSTRAP; + public static final String BROADCAST_FOR_STREAM_ID = STREAM_ID_PREFIX + BROADCAST; + + /* + * Implementation notes: + * Helper for accessing configs related to stream properties. + * + * For most configs, this currently supports two different formats for specifying stream properties: + * 1) "streams.{streamId}.{property}" (recommended to use this format) + * 2) "systems.{systemName}.streams.{streamName}.{property}" (legacy) + * Note that some config lookups are only supported through the "streams.{streamId}.{property}". See the specific + * accessor method to determine which formats are supported. + * + * Summary of terms: + * - streamId: logical identifier used for a stream; configs are specified using this streamId + * - physical stream: concrete name for a stream (if the physical stream is not explicitly configured, then the streamId + * is used as the physical stream + * - streamName: within the javadoc for this class, streamName is the same as physical stream + * - samza property: property which is Samza-specific, which will have "samza." as a prefix (e.g. "samza.key.serde"); + * this is in contrast to stream-specific properties which are related to specific stream technologies + */ + + private Optional nonEmptyOption(String value) { + if (value == null || value.isEmpty()) { + return Optional.empty(); + } else { + return Optional.of(value); + } + } + + public StreamConfig(Config config) { + super(config); + } + + + /** + * Finds the properties from the legacy config style (config key includes system). + * This will return a Config with the properties that match the following formats (if a property is specified through + * multiple formats, priority is top to bottom): + * 1) "systems.{systemName}.streams.{streamName}.{property}" + * 2) "systems.{systemName}.default.stream.{property}" + * + * @param systemName the system name under which the properties are configured + * @param streamName the stream name + * @return the map of properties for the stream + */ + private Map getSystemStreamProperties(String systemName, String streamName) { + if (systemName == null) { + Collections.emptyMap(); + } + SystemConfig systemConfig = new SystemConfig(this); + Config defaults = systemConfig.getDefaultStreamProperties(systemName); + Config explicitConfigs = subset(String.format(STREAM_PREFIX, systemName, streamName), true); + return new MapConfig(defaults, explicitConfigs); + } + + + /** + * Gets all of the properties for the specified streamId (includes current and legacy config styles). + * This will return a Config with the properties that match the following formats (if a property is specified through + * multiple formats, priority is top to bottom): + * 1) "streams.{streamId}.{property}" + * 2) "systems.{systemName}.streams.{streamName}.{property}" where systemName is the system mapped to the streamId in + * the config and streamName is the physical stream name mapped to the stream id + * 3) "systems.{systemName}.default.stream.{property}" where systemName is the system mapped to the streamId in the + * config + * + * @param streamId the identifier for the stream in the config. + * @return the merged map of config properties from both the legacy and new config styles + */ + private MapConfig getAllStreamProperties(String streamId) { + Config allProperties = subset(String.format(STREAM_ID_PREFIX, streamId)); + Map inheritedLegacyProperties = getSystemStreamProperties(getSystem(streamId), getPhysicalName(streamId)); + return new MapConfig(Arrays.asList(inheritedLegacyProperties, allProperties)); + } + + /** + * Gets the distinct stream IDs of all the streams defined in the config + * + * @return collection of stream IDs + */ + public Set getStreamIds() { + // StreamIds are not allowed to have '.' so the first index of '.' marks the end of the streamId. + return subset(STREAMS_PREFIX).keySet().stream().map(key -> key.substring(0, key.indexOf("."))) + .distinct().collect(Collectors.toSet()); + } + + private List getStreamIdsForSystem(String system) { + return getStreamIds().stream().filter(streamId -> system.equals(getSystem(streamId))).collect(Collectors.toList()); + } + + /** + * Finds the stream id which corresponds to the systemStream. + * This finds the stream id that is mapped to the system in systemStream through the config and that has a physical + * name (the physical name might be the stream id itself if there is no explicit mapping) that matches the stream in + * systemStream. + * Note: If the stream in the systemStream is a stream id which is mapped to a physical stream, then that stream won't + * be returned as a stream id here, since the stream in systemStream doesn't match the physical stream name. + * + * @param systemStream system stream to map to stream id + * @return stream id corresponding to the system stream + */ + private String systemStreamToStreamId(SystemStream systemStream) { + List streamIds = getStreamIdsForSystem(systemStream.getSystem()).stream() + .filter(streamId -> systemStream.getStream().equals(getPhysicalName(streamId))).collect(Collectors.toList()); + if (streamIds.size() > 1) { + throw new IllegalStateException(String.format("There was more than one stream found for system stream %s", systemStream)); + } + + return streamIds.isEmpty() ? null : streamIds.get(0); + } + + /** + * Gets the System associated with the specified streamId. + * It first looks for the property + * streams.{streamId}.system + *

+ * If no value was provided, it uses + * job.default.system + * + * @param streamId the identifier for the stream in the config. + * @return the system name associated with the stream or null. + */ + public String getSystem(String streamId) { + String system = get(String.format(SYSTEM_FOR_STREAM_ID, streamId)); + return (system != null) ? system : new JobConfig(this).getDefaultSystem().orElse(null); + } + + /** + * Gets the physical name for the specified streamId. + * + * @param streamId the identifier for the stream in the config. + * @return the physical identifier for the stream or the default if it is undefined. + */ + public String getPhysicalName(String streamId) { + // use streamId as the default physical name + return getOrDefault(String.format(PHYSICAL_NAME_FOR_STREAM_ID, streamId), streamId); + } + + + public Optional getStreamMsgSerde(SystemStream systemStream) { + return nonEmptyOption(getSamzaProperty(systemStream, MSG_SERDE)); + } + + public Optional getStreamKeySerde(SystemStream systemStream) { + return nonEmptyOption(getSamzaProperty(systemStream, KEY_SERDE)); + } + + public boolean getResetOffset(SystemStream systemStream) { + String resetOffset = getSamzaProperty(systemStream, CONSUMER_RESET_OFFSET, "false"); + if (!resetOffset.equalsIgnoreCase("true") && !resetOffset.equalsIgnoreCase("false")) { + LOG.warn("Got a .samza.reset.offset configuration for SystemStream {} that is not true or false (was {})." + + " Defaulting to false.", systemStream, resetOffset); + + resetOffset = "false"; + } + return Boolean.valueOf(resetOffset); + } + + /** + * Determines if a Samza property is specified. + * See getSamzaProperty(SystemStream, String). + * + * @param systemStream the SystemStream for the property value to check + * @param property the samza property key (including the "samza." prefix); for example, for both + * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this + * argument should have the value "samza.prop.key" + */ + protected boolean containsSamzaProperty(SystemStream systemStream, String property) { + if (!property.startsWith(SAMZA_PROPERTY)) { + throw new IllegalArgumentException( + String.format("Attempt to fetch a non samza property for SystemStream %s named %s", systemStream, property)); + } + return getSamzaProperty(systemStream, property) != null; + } + + public boolean isResetOffsetConfigured(SystemStream systemStream) { + return containsSamzaProperty(systemStream, CONSUMER_RESET_OFFSET); + } + + public Optional getDefaultStreamOffset(SystemStream systemStream) { + return Optional.ofNullable(getSamzaProperty(systemStream, CONSUMER_OFFSET_DEFAULT)); + } + + public boolean isDefaultStreamOffsetConfigured(SystemStream systemStream) { + return containsSamzaProperty(systemStream, CONSUMER_OFFSET_DEFAULT); + } + + public boolean getBootstrapEnabled(SystemStream systemStream) { + return Boolean.parseBoolean(getSamzaProperty(systemStream, BOOTSTRAP)); + } + + public boolean getBroadcastEnabled(SystemStream systemStream) { + return Boolean.parseBoolean(getSamzaProperty(systemStream, BROADCAST)); + } + + public int getPriority(SystemStream systemStream) { + return Integer.parseInt(getSamzaProperty(systemStream, PRIORITY, "-1")); + } + + /** + * A streamId is translated to a SystemStream by looking up its System and physicalName. It + * will use the streamId as the stream name if the physicalName doesn't exist. + */ + public SystemStream streamIdToSystemStream(String streamId) { + return new SystemStream(getSystem(streamId), getPhysicalName(streamId)); + } + + + /** + * Returns a list of all SystemStreams that have a serde defined from the config file. + */ + public Set getSerdeStreams(String systemName) { + Config subConf = subset(String.format("systems.%s.streams.", systemName, true)); + Set legacySystemStreams = subConf.keySet().stream() + .filter(k -> k.endsWith(MSG_SERDE) || k.endsWith(KEY_SERDE)) + .map(k -> { + String streamName = k.substring(0, k.length() - 16 /* .samza.XXX.serde length */); + return new SystemStream(systemName, streamName); + }) + .collect(Collectors.toSet()); + + Set systemStreams = subset(STREAMS_PREFIX).keySet().stream() + .filter(k -> k.endsWith(MSG_SERDE) || k.endsWith(KEY_SERDE)) + .map(k -> k.substring(0, k.length() - 16 /* .samza.XXX.serde length */)) + .filter(streamId -> systemName.equals(getSystem(streamId))) + .map(streamId -> streamIdToSystemStream(streamId)).collect(Collectors.toSet()); + + return Sets.union(legacySystemStreams, systemStreams).immutableCopy(); + } + + /* Gets the properties for the streamId which are not Samza properties (i.e. do not have a "samza." prefix). This + * includes current and legacy config styles. + * This will return a Config with the properties that match the following formats (if a property is specified through + * multiple formats, priority is top to bottom): + * 1) "streams.{streamId}.{property}" + * 2) "systems.{systemName}.streams.{streamName}.{property}" where systemName is the system mapped to the streamId in + * the config and streamName is the physical stream name mapped to the stream id + * 3) "systems.{systemName}.default.stream.{property}" where systemName is the system mapped to the streamId in the + * config + * + * @param streamId the identifier for the stream in the config. + * @return the merged map of config properties from both the legacy and new config styles + */ + public Config getStreamProperties(String streamId) { + MapConfig allProperties = getAllStreamProperties(streamId); + Config samzaProperties = allProperties.subset(SAMZA_PROPERTY, false); + Map filteredStreamProperties = + allProperties.entrySet().stream().filter(kv -> !samzaProperties.containsKey(kv.getKey())).collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + return new MapConfig(filteredStreamProperties); + } + + /** + * Gets the boolean flag of whether the specified streamId is an intermediate stream + * + * @param streamId the identifier for the stream in the config. + * @return true if the stream is intermediate + */ + public boolean getIsIntermediateStream(String streamId) { + return getBoolean(String.format(IS_INTERMEDIATE_FOR_STREAM_ID, streamId), false); + } + + /** + * Gets the boolean flag of whether the committed messages specified streamId can be deleted + * + * @param streamId the identifier for the stream in the config. + * @return true if the committed messages of the stream can be deleted + */ + public boolean getDeleteCommittedMessages(String streamId) { + return getBoolean(String.format(DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, streamId), false); + } + + public boolean getIsBounded(String streamId) { + return getBoolean(String.format(IS_BOUNDED_FOR_STREAM_ID, streamId), false); + } + + /** + * Gets the specified Samza property for a SystemStream. A Samza property is a property that controls how Samza + * interacts with the stream, as opposed to a property of the stream itself. + * + * First, tries to map the systemStream to a streamId. This will only find a streamId if the stream is a physical name + * (explicitly mapped physical name or a stream id without a physical name mapping). That means this will not map a + * stream id to itself if there is a mapping from the stream id to a physical stream name. This also requires that the + * stream id is mapped to a system in the config. + * If a stream id is found: + * 1) Look for "streams.{streamId}.{property}" for the stream id. + * 2) Otherwise, look for "systems.{systemName}.streams.{streamName}.{property}" in which the systemName is the system + * mapped to the stream id and the streamName is the physical stream name for the stream id. + * 3) Otherwise, look for "systems.{systemName}.default.stream.{property}" in which the systemName is the system + * mapped to the stream id. + * If a stream id was not found or no property could be found using the above keys: + * 1) Look for "systems.{systemName}.streams.{streamName}.{property}" in which the systemName is the system in the + * input systemStream and the streamName is the stream from the input systemStream. + * 2) Otherwise, look for "systems.{systemName}.default.stream.{property}" in which the systemName is the system + * in the input systemStream. + * 3) or null + * + * @param systemStream the SystemStream for which the property value will be retrieved. + * @param property the samza property key (including the "samza." prefix); for example, for both + * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this + * argument should have the value "samza.prop.key" + * @return the property value + */ + private String getSamzaProperty(SystemStream systemStream, String property) { + if (!property.startsWith(SAMZA_PROPERTY)) { + throw new IllegalArgumentException( + String.format("Attempt to fetch a non samza property for SystemStream %s named %s", systemStream, property)); + } + + String streamVal = getAllStreamProperties(systemStreamToStreamId(systemStream)).get(property); + return (streamVal != null) ? streamVal : getSystemStreamProperties(systemStream.getSystem(), systemStream.getStream()).get(property); + } + + /** + * Gets a Samza property, with a default value used if no property value is found. + * See getSamzaProperty(SystemStream, String). + * + * @param systemStream the SystemStream for which the property value will be retrieved. + * @param property the samza property key (including the "samza." prefix); for example, for both + * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this + * argument should have the value "samza.prop.key" + * @param defaultValue the default value to use if the property value is not found + */ + private String getSamzaProperty(SystemStream systemStream, String property, String defaultValue) { + String streamVal = getSamzaProperty(systemStream, property); + + return streamVal != null ? streamVal : defaultValue; + } +} diff --git a/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.scala b/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.scala deleted file mode 100644 index 6030b93786..0000000000 --- a/samza-core/src/main/scala/org/apache/samza/config/StreamConfig.scala +++ /dev/null @@ -1,381 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you 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 - * - * http://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. - */ - -package org.apache.samza.config - -import org.apache.samza.system.SystemStream -import org.apache.samza.util.Logging - -import scala.collection.JavaConverters._ - -object StreamConfig { - // Samza configs for streams - val SAMZA_PROPERTY = "samza." - val SYSTEM = SAMZA_PROPERTY + "system" - val PHYSICAL_NAME = SAMZA_PROPERTY + "physical.name" - val MSG_SERDE = SAMZA_PROPERTY + "msg.serde" - val KEY_SERDE = SAMZA_PROPERTY + "key.serde" - val CONSUMER_RESET_OFFSET = SAMZA_PROPERTY + "reset.offset" - val CONSUMER_OFFSET_DEFAULT = SAMZA_PROPERTY + "offset.default" - val BOOTSTRAP = SAMZA_PROPERTY + "bootstrap" - val PRIORITY = SAMZA_PROPERTY + "priority" - val IS_INTERMEDIATE = SAMZA_PROPERTY + "intermediate" - val DELETE_COMMITTED_MESSAGES = SAMZA_PROPERTY + "delete.committed.messages" - val IS_BOUNDED = SAMZA_PROPERTY + "bounded" - val BROADCAST = SAMZA_PROPERTY + "broadcast" - - // We don't want any external dependencies on these patterns while both exist. Use getProperty to ensure proper values. - private val STREAMS_PREFIX = "streams." - - val STREAM_PREFIX = "systems.%s.streams.%s." - val STREAM_ID_PREFIX = STREAMS_PREFIX + "%s." - val SYSTEM_FOR_STREAM_ID = STREAM_ID_PREFIX + SYSTEM - val PHYSICAL_NAME_FOR_STREAM_ID = STREAM_ID_PREFIX + PHYSICAL_NAME - val IS_INTERMEDIATE_FOR_STREAM_ID = STREAM_ID_PREFIX + IS_INTERMEDIATE - val DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID = STREAM_ID_PREFIX + DELETE_COMMITTED_MESSAGES - val IS_BOUNDED_FOR_STREAM_ID = STREAM_ID_PREFIX + IS_BOUNDED - val PRIORITY_FOR_STREAM_ID = STREAM_ID_PREFIX + PRIORITY - val CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID = STREAM_ID_PREFIX + CONSUMER_OFFSET_DEFAULT - val BOOTSTRAP_FOR_STREAM_ID = STREAM_ID_PREFIX + BOOTSTRAP - val BROADCAST_FOR_STREAM_ID = STREAM_ID_PREFIX + BROADCAST - - implicit def Config2Stream(config: Config) = new StreamConfig(config) -} - -/** - * Helper for accessing configs related to stream properties. - * - * For most configs, this currently supports two different formats for specifying stream properties: - * 1) "streams.{streamId}.{property}" (recommended to use this format) - * 2) "systems.{systemName}.streams.{streamName}.{property}" (legacy) - * Note that some config lookups are only supported through the "streams.{streamId}.{property}". See the specific - * accessor method to determine which formats are supported. - * - * Summary of terms: - * - streamId: logical identifier used for a stream; configs are specified using this streamId - * - physical stream: concrete name for a stream (if the physical stream is not explicitly configured, then the streamId - * is used as the physical stream - * - streamName: within the javadoc for this class, streamName is the same as physical stream - * - samza property: property which is Samza-specific, which will have "samza." as a prefix (e.g. "samza.key.serde"); - * this is in contrast to stream-specific properties which are related to specific stream technologies - */ -class StreamConfig(config: Config) extends ScalaMapConfig(config) with Logging { - def getStreamMsgSerde(systemStream: SystemStream) = nonEmptyOption(getSamzaProperty(systemStream, StreamConfig.MSG_SERDE)) - - def getStreamKeySerde(systemStream: SystemStream) = nonEmptyOption(getSamzaProperty(systemStream, StreamConfig.KEY_SERDE)) - - def getResetOffset(systemStream: SystemStream) = - Option(getSamzaProperty(systemStream, StreamConfig.CONSUMER_RESET_OFFSET)) match { - case Some("true") => true - case Some("false") => false - case Some(resetOffset) => - warn("Got a .samza.reset.offset configuration for SystemStream %s that is not true or false (was %s). Defaulting to false." - format (systemStream.toString format (systemStream.getSystem, systemStream.getStream), resetOffset)) - false - case _ => false - } - - def isResetOffsetConfigured(systemStream: SystemStream) = - containsSamzaProperty(systemStream, StreamConfig.CONSUMER_RESET_OFFSET) - - def getDefaultStreamOffset(systemStream: SystemStream) = - Option(getSamzaProperty(systemStream, StreamConfig.CONSUMER_OFFSET_DEFAULT)) - - def isDefaultStreamOffsetConfigured(systemStream: SystemStream) = - containsSamzaProperty(systemStream, StreamConfig.CONSUMER_OFFSET_DEFAULT) - - def getBootstrapEnabled(systemStream: SystemStream) = - java.lang.Boolean.parseBoolean(getSamzaProperty(systemStream, StreamConfig.BOOTSTRAP)) - - def getBroadcastEnabled(systemStream: SystemStream) = - java.lang.Boolean.parseBoolean(getSamzaProperty(systemStream, StreamConfig.BROADCAST)) - - def getPriority(systemStream: SystemStream) = - java.lang.Integer.parseInt(getSamzaProperty(systemStream, StreamConfig.PRIORITY, "-1")) - - /** - * Returns a list of all SystemStreams that have a serde defined from the config file. - */ - def getSerdeStreams(systemName: String) = { - val subConf = config.subset("systems.%s.streams." format systemName, true) - val legacySystemStreams = subConf - .asScala - .keys - .filter(k => k.endsWith(StreamConfig.MSG_SERDE) || k.endsWith(StreamConfig.KEY_SERDE)) - .map(k => { - val streamName = k.substring(0, k.length - 16 /* .samza.XXX.serde length */ ) - new SystemStream(systemName, streamName) - }).toSet - - val systemStreams = subset(StreamConfig.STREAMS_PREFIX) - .asScala - .keys - .filter(k => k.endsWith(StreamConfig.MSG_SERDE) || k.endsWith(StreamConfig.KEY_SERDE)) - .map(k => k.substring(0, k.length - 16 /* .samza.XXX.serde length */ )) - .filter(streamId => systemName.equals(getSystem(streamId))) - .map(streamId => streamIdToSystemStream(streamId)).toSet - - legacySystemStreams.union(systemStreams) - } - - /** - * Gets the properties for the streamId which are not Samza properties (i.e. do not have a "samza." prefix). This - * includes current and legacy config styles. - * This will return a Config with the properties that match the following formats (if a property is specified through - * multiple formats, priority is top to bottom): - * 1) "streams.{streamId}.{property}" - * 2) "systems.{systemName}.streams.{streamName}.{property}" where systemName is the system mapped to the streamId in - * the config and streamName is the physical stream name mapped to the stream id - * 3) "systems.{systemName}.default.stream.{property}" where systemName is the system mapped to the streamId in the - * config - * - * @param streamId the identifier for the stream in the config. - * @return the merged map of config properties from both the legacy and new config styles - */ - def getStreamProperties(streamId: String) = { - val allProperties = getAllStreamProperties(streamId) - val samzaProperties = allProperties.subset(StreamConfig.SAMZA_PROPERTY, false) - val filteredStreamProperties: java.util.Map[String, String] = - allProperties.asScala.filterKeys(k => !samzaProperties.containsKey(k)).asJava - new MapConfig(filteredStreamProperties) - } - - /** - * Gets the System associated with the specified streamId. - * It first looks for the property - * streams.{streamId}.system - * - * If no value was provided, it uses - * job.default.system - * - * @param streamId the identifier for the stream in the config. - * @return the system name associated with the stream. - */ - def getSystem(streamId: String) = { - getOption(StreamConfig.SYSTEM_FOR_STREAM_ID format streamId) match { - case Some(system) => system - case _ => new JobConfig(config).getDefaultSystem().orElse(null) - } - } - - /** - * Gets the physical name for the specified streamId. - * - * @param streamId the identifier for the stream in the config. - * @return the physical identifier for the stream or the default if it is undefined. - */ - def getPhysicalName(streamId: String) = { - // use streamId as the default physical name - get(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID format streamId, streamId) - } - - /** - * Gets the boolean flag of whether the specified streamId is an intermediate stream - * @param streamId the identifier for the stream in the config. - * @return true if the stream is intermediate - */ - def getIsIntermediateStream(streamId: String) = { - getBoolean(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID format streamId, false) - } - - /** - * Gets the boolean flag of whether the committed messages specified streamId can be deleted - * @param streamId the identifier for the stream in the config. - * @return true if the committed messages of the stream can be deleted - */ - def getDeleteCommittedMessages(streamId: String) = { - getBoolean(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID format streamId, false) - } - - def getIsBounded(streamId: String) = { - getBoolean(StreamConfig.IS_BOUNDED_FOR_STREAM_ID format streamId, false) - } - - /** - * Gets the stream IDs of all the streams defined in the config - * @return collection of stream IDs - */ - def getStreamIds(): Iterable[String] = { - // StreamIds are not allowed to have '.' so the first index of '.' marks the end of the streamId. - subset(StreamConfig.STREAMS_PREFIX).asScala.keys.map(key => key.substring(0, key.indexOf("."))) - } - - /** - * Gets the specified Samza property for a SystemStream. A Samza property is a property that controls how Samza - * interacts with the stream, as opposed to a property of the stream itself. - * - * First, tries to map the systemStream to a streamId. This will only find a streamId if the stream is a physical name - * (explicitly mapped physical name or a stream id without a physical name mapping). That means this will not map a - * stream id to itself if there is a mapping from the stream id to a physical stream name. This also requires that the - * stream id is mapped to a system in the config. - * If a stream id is found: - * 1) Look for "streams.{streamId}.{property}" for the stream id. - * 2) Otherwise, look for "systems.{systemName}.streams.{streamName}.{property}" in which the systemName is the system - * mapped to the stream id and the streamName is the physical stream name for the stream id. - * 3) Otherwise, look for "systems.{systemName}.default.stream.{property}" in which the systemName is the system - * mapped to the stream id. - * If a stream id was not found or no property could be found using the above keys: - * 1) Look for "systems.{systemName}.streams.{streamName}.{property}" in which the systemName is the system in the - * input systemStream and the streamName is the stream from the input systemStream. - * 2) Otherwise, look for "systems.{systemName}.default.stream.{property}" in which the systemName is the system - * in the input systemStream. - * - * @param systemStream the SystemStream for which the property value will be retrieved. - * @param property the samza property key (including the "samza." prefix); for example, for both - * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this - * argument should have the value "samza.prop.key" - */ - private def getSamzaProperty(systemStream: SystemStream, property: String): String = { - if (!property.startsWith(StreamConfig.SAMZA_PROPERTY)) { - throw new IllegalArgumentException( - "Attempt to fetch a non samza property for SystemStream %s named %s" format(systemStream, property)) - } - - val streamVal = getAllStreamProperties(systemStreamToStreamId(systemStream)).get(property) - if (streamVal != null) { - streamVal - } else { - getSystemStreamProperties(systemStream.getSystem(), systemStream.getStream).get(property) - } - } - - /** - * Gets a Samza property, with a default value used if no property value is found. - * See getSamzaProperty(SystemStream, String). - * - * @param systemStream the SystemStream for which the property value will be retrieved. - * @param property the samza property key (including the "samza." prefix); for example, for both - * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this - * argument should have the value "samza.prop.key" - * @param defaultValue the default value to use if the property value is not found - */ - private def getSamzaProperty(systemStream: SystemStream, property: String, defaultValue: String): String = { - val streamVal = getSamzaProperty(systemStream, property) - - if (streamVal != null) { - streamVal - } else { - defaultValue - } - } - - /** - * Determines if a Samza property is specified. - * See getSamzaProperty(SystemStream, String). - * - * @param systemStream the SystemStream for the property value to check - * @param property the samza property key (including the "samza." prefix); for example, for both - * "streams.streamId.samza.prop.key" and "systems.system.streams.streamName.samza.prop.key", this - * argument should have the value "samza.prop.key" - */ - private def containsSamzaProperty(systemStream: SystemStream, property: String): Boolean = { - if (!property.startsWith(StreamConfig.SAMZA_PROPERTY)) { - throw new IllegalArgumentException( - "Attempt to fetch a non samza property for SystemStream %s named %s" format(systemStream, property)) - } - getSamzaProperty(systemStream, property) != null - } - - - /** - * Finds the properties from the legacy config style (config key includes system). - * This will return a Config with the properties that match the following formats (if a property is specified through - * multiple formats, priority is top to bottom): - * 1) "systems.{systemName}.streams.{streamName}.{property}" - * 2) "systems.{systemName}.default.stream.{property}" - * - * @param systemName the system name under which the properties are configured - * @param streamName the stream name - * @return the map of properties for the stream - */ - private def getSystemStreamProperties(systemName: String, streamName: String) = { - if (systemName == null) { - Map() - } - val systemConfig = new SystemConfig(config) - val defaults = systemConfig.getDefaultStreamProperties(systemName) - val explicitConfigs = config.subset(StreamConfig.STREAM_PREFIX format(systemName, streamName), true) - new MapConfig(defaults, explicitConfigs) - } - - /** - * Gets all of the properties for the specified streamId (includes current and legacy config styles). - * This will return a Config with the properties that match the following formats (if a property is specified through - * multiple formats, priority is top to bottom): - * 1) "streams.{streamId}.{property}" - * 2) "systems.{systemName}.streams.{streamName}.{property}" where systemName is the system mapped to the streamId in - * the config and streamName is the physical stream name mapped to the stream id - * 3) "systems.{systemName}.default.stream.{property}" where systemName is the system mapped to the streamId in the - * config - * - * @param streamId the identifier for the stream in the config. - * @return the merged map of config properties from both the legacy and new config styles - */ - private def getAllStreamProperties(streamId: String) = { - val allProperties = subset(StreamConfig.STREAM_ID_PREFIX format streamId) - val inheritedLegacyProperties: java.util.Map[String, String] = - getSystemStreamProperties(getSystem(streamId), getPhysicalName(streamId)) - new MapConfig(java.util.Arrays.asList(inheritedLegacyProperties, allProperties)) - } - - private def getStreamIdsForSystem(system: String): Iterable[String] = { - getStreamIds().filter(streamId => system.equals(getSystem(streamId))) - } - - /** - * Finds the stream id which corresponds to the systemStream. - * This finds the stream id that is mapped to the system in systemStream through the config and that has a physical - * name (the physical name might be the stream id itself if there is no explicit mapping) that matches the stream in - * systemStream. - * Note: If the stream in the systemStream is a stream id which is mapped to a physical stream, then that stream won't - * be returned as a stream id here, since the stream in systemStream doesn't match the physical stream name. - * - * @param systemStream system stream to map to stream id - * @return stream id corresponding to the system stream - */ - private def systemStreamToStreamId(systemStream: SystemStream): String = { - val streamIds = getStreamIdsForSystem(systemStream.getSystem) - .filter(streamId => systemStream.getStream().equals(getPhysicalName(streamId))) - if (streamIds.size > 1) { - throw new IllegalStateException("There was more than one stream found for system stream %s" format(systemStream)) - } - - if (streamIds.isEmpty) { - null - } else { - streamIds.head - } - } - - /** - * A streamId is translated to a SystemStream by looking up its System and physicalName. It - * will use the streamId as the stream name if the physicalName doesn't exist. - */ - def streamIdToSystemStream(streamId: String): SystemStream = { - new SystemStream(getSystem(streamId), getPhysicalName(streamId)) - } - - private def nonEmptyOption(value: String): Option[String] = { - if (value == null || value.isEmpty) { - None - } else { - Some(value) - } - } -} diff --git a/samza-core/src/main/scala/org/apache/samza/container/SamzaContainer.scala b/samza-core/src/main/scala/org/apache/samza/container/SamzaContainer.scala index f7c75f4d10..b552860225 100644 --- a/samza-core/src/main/scala/org/apache/samza/container/SamzaContainer.scala +++ b/samza-core/src/main/scala/org/apache/samza/container/SamzaContainer.scala @@ -25,14 +25,14 @@ import java.net.{URL, UnknownHostException} import java.nio.file.Path import java.time.Duration import java.util -import java.util.Base64 +import java.util.{Base64, Optional} import java.util.concurrent.{ExecutorService, Executors, ScheduledExecutorService, TimeUnit} +import java.util.stream.Collectors import com.google.common.annotations.VisibleForTesting import com.google.common.util.concurrent.ThreadFactoryBuilder import org.apache.samza.checkpoint.{CheckpointListener, OffsetManager, OffsetManagerMetrics} -import org.apache.samza.config.StreamConfig.Config2Stream -import org.apache.samza.config._ +import org.apache.samza.config.{StreamConfig, _} import org.apache.samza.container.disk.DiskSpaceMonitor.Listener import org.apache.samza.container.disk.{DiskQuotaPolicyFactory, DiskSpaceMonitor, NoThrottlingDiskQuotaPolicyFactory, PollingScanDiskSpaceMonitor} import org.apache.samza.container.host.{StatisticsMonitorImpl, SystemMemoryStatistics, SystemStatisticsMonitor} @@ -196,7 +196,8 @@ object SamzaContainer extends Logging { info("Got system names: %s" format systemNames) - val serdeStreams = systemNames.foldLeft(Set[SystemStream]())(_ ++ config.getSerdeStreams(_)) + val streamConfig = new StreamConfig(config) + val serdeStreams = systemNames.foldLeft(Set[SystemStream]())(_ ++ streamConfig.getSerdeStreams(_).asScala) info("Got serde streams: %s" format serdeStreams) @@ -302,9 +303,9 @@ object SamzaContainer extends Logging { * A Helper function to build a Map[SystemStream, Serde] for streams defined in the config. * This is useful to build both key and message serde maps. */ - val buildSystemStreamSerdeMap = (getSerdeName: (SystemStream) => Option[String]) => { + val buildSystemStreamSerdeMap = (getSerdeName: (SystemStream) => Optional[String]) => { (serdeStreams ++ inputSystemStreamPartitions) - .filter(systemStream => getSerdeName(systemStream).isDefined) + .filter(systemStream => getSerdeName(systemStream).isPresent) .flatMap(systemStream => { val serdeName = getSerdeName(systemStream).get val serde = serdes.getOrElse(serdeName, @@ -327,11 +328,11 @@ object SamzaContainer extends Logging { debug("Got system message serdes: %s" format systemMessageSerdes) - val systemStreamKeySerdes = buildSystemStreamSerdeMap(systemStream => config.getStreamKeySerde(systemStream)) + val systemStreamKeySerdes = buildSystemStreamSerdeMap(systemStream => streamConfig.getStreamKeySerde(systemStream)) debug("Got system stream key serdes: %s" format systemStreamKeySerdes) - val systemStreamMessageSerdes = buildSystemStreamSerdeMap(systemStream => config.getStreamMsgSerde(systemStream)) + val systemStreamMessageSerdes = buildSystemStreamSerdeMap(systemStream => streamConfig.getStreamMsgSerde(systemStream)) debug("Got system stream message serdes: %s" format systemStreamMessageSerdes) @@ -360,16 +361,17 @@ object SamzaContainer extends Logging { SystemClock.instance, getChangelogSSPsForContainer(containerModel, changeLogSystemStreams).asJava) - val intermediateStreams = config - .getStreamIds - .filter(config.getIsIntermediateStream(_)) + val intermediateStreams = streamConfig + .getStreamIds() + .asScala + .filter((streamId:String) => streamConfig.getIsIntermediateStream(streamId)) .toList info("Got intermediate streams: %s" format intermediateStreams) val controlMessageKeySerdes = intermediateStreams .flatMap(streamId => { - val systemStream = config.streamIdToSystemStream(streamId) + val systemStream = streamConfig.streamIdToSystemStream(streamId) systemStreamKeySerdes.get(systemStream) .orElse(systemKeySerdes.get(systemStream.getSystem)) .map(serde => (systemStream, new StringSerde("UTF-8"))) @@ -377,7 +379,7 @@ object SamzaContainer extends Logging { val intermediateStreamMessageSerdes = intermediateStreams .flatMap(streamId => { - val systemStream = config.streamIdToSystemStream(streamId) + val systemStream = streamConfig.streamIdToSystemStream(streamId) systemStreamMessageSerdes.get(systemStream) .orElse(systemMessageSerdes.get(systemStream.getSystem)) .map(serde => (systemStream, new IntermediateMessageSerde(serde))) diff --git a/samza-core/src/main/scala/org/apache/samza/container/TaskInstance.scala b/samza-core/src/main/scala/org/apache/samza/container/TaskInstance.scala index a17a790156..f66f4c7bb2 100644 --- a/samza-core/src/main/scala/org/apache/samza/container/TaskInstance.scala +++ b/samza-core/src/main/scala/org/apache/samza/container/TaskInstance.scala @@ -22,10 +22,10 @@ package org.apache.samza.container import java.util.{Objects, Optional} import java.util.concurrent.ScheduledExecutorService + import org.apache.samza.SamzaException import org.apache.samza.checkpoint.OffsetManager -import org.apache.samza.config.Config -import org.apache.samza.config.StreamConfig.Config2Stream +import org.apache.samza.config.{Config, StreamConfig} import org.apache.samza.context._ import org.apache.samza.job.model.{JobModel, TaskModel} import org.apache.samza.scheduler.{CallbackSchedulerImpl, EpochTimeScheduler, ScheduledCallback} @@ -99,9 +99,10 @@ class TaskInstance( private val config: Config = jobContext.getConfig - val intermediateStreams: Set[String] = config.getStreamIds.filter(config.getIsIntermediateStream).toSet + val streamConfig: StreamConfig = new StreamConfig(config) + val intermediateStreams: Set[String] = streamConfig.getStreamIds.filter(streamConfig.getIsIntermediateStream).toSet - val streamsToDeleteCommittedMessages: Set[String] = config.getStreamIds.filter(config.getDeleteCommittedMessages).map(config.getPhysicalName).toSet + val streamsToDeleteCommittedMessages: Set[String] = streamConfig.getStreamIds.filter(streamConfig.getDeleteCommittedMessages).map(streamConfig.getPhysicalName).toSet def registerOffsets { debug("Registering offsets for taskName: %s" format taskName) diff --git a/samza-core/src/main/scala/org/apache/samza/metrics/reporter/MetricsSnapshotReporterFactory.scala b/samza-core/src/main/scala/org/apache/samza/metrics/reporter/MetricsSnapshotReporterFactory.scala index 27201c952a..ec5f4d971f 100644 --- a/samza-core/src/main/scala/org/apache/samza/metrics/reporter/MetricsSnapshotReporterFactory.scala +++ b/samza-core/src/main/scala/org/apache/samza/metrics/reporter/MetricsSnapshotReporterFactory.scala @@ -21,8 +21,7 @@ package org.apache.samza.metrics.reporter import org.apache.samza.util.{Logging, StreamUtil, Util} import org.apache.samza.SamzaException -import org.apache.samza.config.{Config, JobConfig, MetricsConfig, SerializerConfig, SystemConfig} -import org.apache.samza.config.StreamConfig.Config2Stream +import org.apache.samza.config.{Config, JobConfig, MetricsConfig, SerializerConfig, StreamConfig, SystemConfig} import org.apache.samza.metrics.MetricsReporter import org.apache.samza.metrics.MetricsReporterFactory import org.apache.samza.metrics.MetricsRegistryMap @@ -63,10 +62,11 @@ class MetricsSnapshotReporterFactory extends MetricsReporterFactory with Logging val producer = systemFactory.getProducer(systemName, config, registry) info("Got producer %s." format producer) + val streamConfig = new StreamConfig(config) - val streamSerdeName = config.getStreamMsgSerde(systemStream) - val systemSerdeName = JavaOptionals.toRichOptional(systemConfig.getSystemMsgSerde(systemName)).toOption - val serdeName = streamSerdeName.getOrElse(systemSerdeName.getOrElse(null)) + val streamSerdeName = streamConfig.getStreamMsgSerde(systemStream) + val systemSerdeName = systemConfig.getSystemMsgSerde(systemName) + val serdeName = streamSerdeName.orElse(systemSerdeName.orElse(null)) val serializerConfig = new SerializerConfig(config) val serde = if (serdeName != null) { JavaOptionals.toRichOptional(serializerConfig.getSerdeFactoryClass(serdeName)).toOption match { diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStreamConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStreamConfig.java index b91eeb0c19..bbea19e528 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStreamConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStreamConfig.java @@ -18,14 +18,15 @@ */ package org.apache.samza.config; -import java.util.Collections; import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; import com.google.common.collect.Sets; import org.apache.samza.system.SystemStream; import org.junit.Test; -import scala.Option; -import scala.collection.JavaConverters; + +import java.util.Collections; +import java.util.Optional; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -47,51 +48,51 @@ public class TestStreamConfig { @Test public void testGetStreamMsgSerde() { String value = "my.msg.serde"; - doTestSamzaProperty(StreamConfig.MSG_SERDE(), value, - (config, systemStream) -> assertEquals(Option.apply(value), config.getStreamMsgSerde(systemStream))); - doTestSamzaProperty(StreamConfig.MSG_SERDE(), "", - (config, systemStream) -> assertEquals(Option.empty(), config.getStreamMsgSerde(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.MSG_SERDE(), - (config, systemStream) -> assertEquals(Option.empty(), config.getStreamMsgSerde(systemStream))); + doTestSamzaProperty(StreamConfig.MSG_SERDE, value, + (config, systemStream) -> assertEquals(Optional.of(value), config.getStreamMsgSerde(systemStream))); + doTestSamzaProperty(StreamConfig.MSG_SERDE, "", + (config, systemStream) -> assertEquals(Optional.empty(), config.getStreamMsgSerde(systemStream))); + doTestSamzaPropertyDoesNotExist(StreamConfig.MSG_SERDE, + (config, systemStream) -> assertEquals(Optional.empty(), config.getStreamMsgSerde(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getStreamMsgSerde); } @Test public void testGetStreamKeySerde() { String value = "my.key.serde"; - doTestSamzaProperty(StreamConfig.KEY_SERDE(), value, - (config, systemStream) -> assertEquals(Option.apply(value), config.getStreamKeySerde(systemStream))); - doTestSamzaProperty(StreamConfig.KEY_SERDE(), "", - (config, systemStream) -> assertEquals(Option.empty(), config.getStreamKeySerde(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.KEY_SERDE(), - (config, systemStream) -> assertEquals(Option.empty(), config.getStreamKeySerde(systemStream))); + doTestSamzaProperty(StreamConfig.KEY_SERDE, value, + (config, systemStream) -> assertEquals(Optional.of(value), config.getStreamKeySerde(systemStream))); + doTestSamzaProperty(StreamConfig.KEY_SERDE, "", + (config, systemStream) -> assertEquals(Optional.empty(), config.getStreamKeySerde(systemStream))); + doTestSamzaPropertyDoesNotExist(StreamConfig.KEY_SERDE, + (config, systemStream) -> assertEquals(Optional.empty(), config.getStreamKeySerde(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getStreamKeySerde); } @Test public void testGetResetOffset() { - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "true", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "true", (config, systemStream) -> assertTrue(config.getResetOffset(systemStream))); - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "false", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "false", (config, systemStream) -> assertFalse(config.getResetOffset(systemStream))); // if not true/false, then use false - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "unknown_value", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "unknown_value", (config, systemStream) -> assertFalse(config.getResetOffset(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_RESET_OFFSET(), + doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_RESET_OFFSET, (config, systemStream) -> assertFalse(config.getResetOffset(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getResetOffset); } @Test public void testIsResetOffsetConfigured() { - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "true", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "true", (config, systemStream) -> assertTrue(config.isResetOffsetConfigured(systemStream))); - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "false", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "false", (config, systemStream) -> assertTrue(config.isResetOffsetConfigured(systemStream))); // if not true/false, then use false - doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET(), "unknown_value", + doTestSamzaProperty(StreamConfig.CONSUMER_RESET_OFFSET, "unknown_value", (config, systemStream) -> assertTrue(config.isResetOffsetConfigured(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_RESET_OFFSET(), + doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_RESET_OFFSET, (config, systemStream) -> assertFalse(config.isResetOffsetConfigured(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::isResetOffsetConfigured); } @@ -99,12 +100,12 @@ public void testIsResetOffsetConfigured() { @Test public void testGetDefaultStreamOffset() { String value = "my_offset_default"; - doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT(), value, - (config, systemStream) -> assertEquals(Option.apply(value), config.getDefaultStreamOffset(systemStream))); - doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT(), "", - (config, systemStream) -> assertEquals(Option.apply(""), config.getDefaultStreamOffset(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_OFFSET_DEFAULT(), - (config, systemStream) -> assertEquals(Option.empty(), + doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT, value, + (config, systemStream) -> assertEquals(Optional.of(value), config.getDefaultStreamOffset(systemStream))); + doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT, "", + (config, systemStream) -> assertEquals(Optional.of(""), config.getDefaultStreamOffset(systemStream))); + doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_OFFSET_DEFAULT, + (config, systemStream) -> assertEquals(Optional.empty(), new StreamConfig(config).getDefaultStreamOffset(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getDefaultStreamOffset); } @@ -112,52 +113,52 @@ public void testGetDefaultStreamOffset() { @Test public void testIsDefaultStreamOffsetConfigured() { String value = "my_offset_default"; - doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT(), value, + doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT, value, (config, systemStream) -> assertTrue(config.isDefaultStreamOffsetConfigured(systemStream))); - doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT(), "", + doTestSamzaProperty(StreamConfig.CONSUMER_OFFSET_DEFAULT, "", (config, systemStream) -> assertTrue(config.isDefaultStreamOffsetConfigured(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_OFFSET_DEFAULT(), + doTestSamzaPropertyDoesNotExist(StreamConfig.CONSUMER_OFFSET_DEFAULT, (config, systemStream) -> assertFalse(config.isDefaultStreamOffsetConfigured(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::isDefaultStreamOffsetConfigured); } @Test public void testGetBootstrapEnabled() { - doTestSamzaProperty(StreamConfig.BOOTSTRAP(), "true", + doTestSamzaProperty(StreamConfig.BOOTSTRAP, "true", (config, systemStream) -> assertTrue(config.getBootstrapEnabled(systemStream))); - doTestSamzaProperty(StreamConfig.BOOTSTRAP(), "false", + doTestSamzaProperty(StreamConfig.BOOTSTRAP, "false", (config, systemStream) -> assertFalse(config.getBootstrapEnabled(systemStream))); // if not true/false, then use false - doTestSamzaProperty(StreamConfig.BOOTSTRAP(), "unknown_value", + doTestSamzaProperty(StreamConfig.BOOTSTRAP, "unknown_value", (config, systemStream) -> assertFalse(config.getBootstrapEnabled(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.BOOTSTRAP(), + doTestSamzaPropertyDoesNotExist(StreamConfig.BOOTSTRAP, (config, systemStream) -> assertFalse(config.getBootstrapEnabled(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getBootstrapEnabled); } @Test public void testGetBroadcastEnabled() { - doTestSamzaProperty(StreamConfig.BROADCAST(), "true", + doTestSamzaProperty(StreamConfig.BROADCAST, "true", (config, systemStream) -> assertTrue(config.getBroadcastEnabled(systemStream))); - doTestSamzaProperty(StreamConfig.BROADCAST(), "false", + doTestSamzaProperty(StreamConfig.BROADCAST, "false", (config, systemStream) -> assertFalse(config.getBroadcastEnabled(systemStream))); // if not true/false, then use false - doTestSamzaProperty(StreamConfig.BROADCAST(), "unknown_value", + doTestSamzaProperty(StreamConfig.BROADCAST, "unknown_value", (config, systemStream) -> assertFalse(config.getBroadcastEnabled(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.BROADCAST(), + doTestSamzaPropertyDoesNotExist(StreamConfig.BROADCAST, (config, systemStream) -> assertFalse(config.getBroadcastEnabled(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getBroadcastEnabled); } @Test public void testGetPriority() { - doTestSamzaProperty(StreamConfig.PRIORITY(), "0", + doTestSamzaProperty(StreamConfig.PRIORITY, "0", (config, systemStream) -> assertEquals(0, config.getPriority(systemStream))); - doTestSamzaProperty(StreamConfig.PRIORITY(), "100", + doTestSamzaProperty(StreamConfig.PRIORITY, "100", (config, systemStream) -> assertEquals(100, config.getPriority(systemStream))); - doTestSamzaProperty(StreamConfig.PRIORITY(), "-1", + doTestSamzaProperty(StreamConfig.PRIORITY, "-1", (config, systemStream) -> assertEquals(-1, config.getPriority(systemStream))); - doTestSamzaPropertyDoesNotExist(StreamConfig.PRIORITY(), + doTestSamzaPropertyDoesNotExist(StreamConfig.PRIORITY, (config, systemStream) -> assertEquals(-1, config.getPriority(systemStream))); doTestSamzaPropertyInvalidConfig(StreamConfig::getPriority); } @@ -165,78 +166,68 @@ public void testGetPriority() { @Test public void testGetSerdeStreams() { assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(new StreamConfig(new MapConfig()).getSerdeStreams(SYSTEM)).asJava()); + new StreamConfig(new MapConfig()).getSerdeStreams(SYSTEM)); // not key/msg serde property for "streams." StreamConfig streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); - assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); + assertEquals(Collections.emptySet(), streamConfig.getSerdeStreams(SYSTEM)); // not matching system for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.KEY_SERDE(), UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), "otherSystem"))); - assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.KEY_SERDE, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), "otherSystem"))); + assertEquals(Collections.emptySet(), streamConfig.getSerdeStreams(SYSTEM)); // not key/msg serde property for "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); - assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); + assertEquals(Collections.emptySet(), streamConfig.getSerdeStreams(SYSTEM)); // not matching system for "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), "otherSystem", STREAM_ID) + StreamConfig.KEY_SERDE(), + String.format(StreamConfig.STREAM_PREFIX, "otherSystem", STREAM_ID) + StreamConfig.KEY_SERDE, UNUSED_VALUE))); - assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + assertEquals(Collections.emptySet(), streamConfig.getSerdeStreams(SYSTEM)); // not matching system for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.KEY_SERDE(), UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), "otherSystem"))); - assertEquals(Collections.emptySet(), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.KEY_SERDE, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), "otherSystem"))); + assertEquals(Collections.emptySet(), streamConfig.getSerdeStreams(SYSTEM)); String serdeValue = "my.serde.class"; // key serde for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.KEY_SERDE(), serdeValue))); - assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.KEY_SERDE, serdeValue))); + assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), streamConfig.getSerdeStreams(SYSTEM)); // msg serde for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.MSG_SERDE(), serdeValue))); - assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.MSG_SERDE, serdeValue))); + assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), streamConfig.getSerdeStreams(SYSTEM)); // serde for "streams." with physical stream name mapping streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM, - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.KEY_SERDE(), serdeValue))); - assertEquals(Collections.singleton(new SystemStream(SYSTEM, PHYSICAL_STREAM)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.KEY_SERDE, serdeValue))); + assertEquals(Collections.singleton(new SystemStream(SYSTEM, PHYSICAL_STREAM)), streamConfig.getSerdeStreams(SYSTEM)); // key serde for "systems..streams." streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + StreamConfig.KEY_SERDE(), + ImmutableMap.of(String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + StreamConfig.KEY_SERDE, serdeValue))); - assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), streamConfig.getSerdeStreams(SYSTEM)); // msg serde for "systems..streams." streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + StreamConfig.MSG_SERDE(), + ImmutableMap.of(String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + StreamConfig.MSG_SERDE, serdeValue))); - assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + assertEquals(Collections.singleton(new SystemStream(SYSTEM, STREAM_ID)), streamConfig.getSerdeStreams(SYSTEM)); // merge several different ways of providing serdes String streamIdWithPhysicalName = "streamIdWithPhysicalName"; @@ -244,19 +235,18 @@ public void testGetSerdeStreams() { // need to map the stream ids to the system .put(JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM) // key and msg serde for "streams." - .put(String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.KEY_SERDE(), serdeValue) - .put(String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + StreamConfig.MSG_SERDE(), serdeValue) + .put(String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.KEY_SERDE, serdeValue) + .put(String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + StreamConfig.MSG_SERDE, serdeValue) // key serde for "streams." with physical stream name mapping - .put(String.format(StreamConfig.STREAM_ID_PREFIX(), streamIdWithPhysicalName) + StreamConfig.KEY_SERDE(), + .put(String.format(StreamConfig.STREAM_ID_PREFIX, streamIdWithPhysicalName) + StreamConfig.KEY_SERDE, serdeValue) - .put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), streamIdWithPhysicalName), PHYSICAL_STREAM) + .put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, streamIdWithPhysicalName), PHYSICAL_STREAM) // key serde for "systems..streams." - .put(String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, OTHER_STREAM_ID) + StreamConfig.KEY_SERDE(), + .put(String.format(StreamConfig.STREAM_PREFIX, SYSTEM, OTHER_STREAM_ID) + StreamConfig.KEY_SERDE, serdeValue) .build())); assertEquals(Sets.newHashSet(new SystemStream(SYSTEM, STREAM_ID), new SystemStream(SYSTEM, PHYSICAL_STREAM), - new SystemStream(SYSTEM, OTHER_STREAM_ID)), - JavaConverters.setAsJavaSetConverter(streamConfig.getSerdeStreams(SYSTEM)).asJava()); + new SystemStream(SYSTEM, OTHER_STREAM_ID)), streamConfig.getSerdeStreams(SYSTEM)); } @Test @@ -269,36 +259,36 @@ public void testGetStreamProperties() { // not matching stream id for "streams." StreamConfig streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), OTHER_STREAM_ID) + propertyName, UNUSED_VALUE))); + String.format(StreamConfig.STREAM_ID_PREFIX, OTHER_STREAM_ID) + propertyName, UNUSED_VALUE))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(STREAM_ID)); // not matching stream id for "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, OTHER_STREAM_ID) + propertyName, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), OTHER_STREAM_ID), SYSTEM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, OTHER_STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, OTHER_STREAM_ID), SYSTEM))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(STREAM_ID)); // no system mapping when using "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(STREAM_ID)); // ignore property with "samza" prefix for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(STREAM_ID)); // ignore property with "samza" prefix for "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(STREAM_ID)); // should not map physical name back to stream id if physical name is passed as stream id streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertEquals(new MapConfig(), streamConfig.getStreamProperties(PHYSICAL_STREAM)); // BEGIN: tests in which properties can be found in the config @@ -307,64 +297,64 @@ public void testGetStreamProperties() { // "streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, propertyValue))); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, propertyValue))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, propertyValue, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, propertyValue, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // "systems..default.stream." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, propertyValue, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // use physical name mapping for "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, PHYSICAL_STREAM) + propertyName, propertyValue, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, PHYSICAL_STREAM) + propertyName, propertyValue, // should not use stream id since there is physical stream - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // "streams." should override "systems..streams." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, propertyValue, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, propertyValue, // should not use "systems..streams." since there is a "streams." config - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // "systems..streams." should override "systems..default.stream." streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, propertyValue, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, propertyValue, // should not use "systems..default.stream." since there is a "systems..streams." String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, UNUSED_VALUE, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(ImmutableMap.of(propertyName, propertyValue)), streamConfig.getStreamProperties(STREAM_ID)); // merge multiple ways of specifying configs streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( // "streams." - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + "from.stream.id.property", "fromStreamIdValue", + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + "from.stream.id.property", "fromStreamIdValue", // second "streams." property - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + "from.stream.id.other.property", + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + "from.stream.id.other.property", "fromStreamIdOtherValue", // "systems..streams." - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + "from.system.stream.property", + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + "from.system.stream.property", "fromSystemStreamValue", // need to map the stream id to a system - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))); assertEquals(new MapConfig(ImmutableMap.of( "from.stream.id.property", "fromStreamIdValue", "from.stream.id.other.property", "fromStreamIdOtherValue", @@ -378,15 +368,15 @@ public void testGetSystem() { // system is specified directly StreamConfig streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, JobConfig.JOB_DEFAULT_SYSTEM, "otherSystem", - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), OTHER_STREAM_ID), "otherSystem"))); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, OTHER_STREAM_ID), "otherSystem"))); assertEquals(SYSTEM, streamConfig.getSystem(STREAM_ID)); // fall back to job default system streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), OTHER_STREAM_ID), "otherSystem"))); + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, OTHER_STREAM_ID), "otherSystem"))); assertEquals(SYSTEM, streamConfig.getSystem(STREAM_ID)); } @@ -396,11 +386,11 @@ public void testGetPhysicalName() { // ignore mapping for other stream ids StreamConfig streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), OTHER_STREAM_ID), PHYSICAL_STREAM))); + ImmutableMap.of(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, OTHER_STREAM_ID), PHYSICAL_STREAM))); assertEquals(STREAM_ID, streamConfig.getPhysicalName(STREAM_ID)); streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + ImmutableMap.of(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertEquals(PHYSICAL_STREAM, streamConfig.getPhysicalName(STREAM_ID)); } @@ -410,21 +400,21 @@ public void testGetIsIntermediateStream() { // ignore mapping for other stream ids StreamConfig streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID(), OTHER_STREAM_ID), "true"))); + ImmutableMap.of(String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID, OTHER_STREAM_ID), "true"))); assertFalse(streamConfig.getIsIntermediateStream(STREAM_ID)); // do not use stream id property if physical name is passed as input streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID(), STREAM_ID), "true", - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID, STREAM_ID), "true", + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertFalse(streamConfig.getIsIntermediateStream(PHYSICAL_STREAM)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID(), STREAM_ID), "true"))); + String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID, STREAM_ID), "true"))); assertTrue(streamConfig.getIsIntermediateStream(STREAM_ID)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID(), STREAM_ID), "false"))); + String.format(StreamConfig.IS_INTERMEDIATE_FOR_STREAM_ID, STREAM_ID), "false"))); assertFalse(streamConfig.getIsIntermediateStream(STREAM_ID)); } @@ -434,22 +424,22 @@ public void testGetDeleteCommittedMessages() { // ignore mapping for other stream ids StreamConfig streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID(), OTHER_STREAM_ID), + ImmutableMap.of(String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, OTHER_STREAM_ID), "true"))); assertFalse(streamConfig.getDeleteCommittedMessages(STREAM_ID)); // do not use stream id property if physical name is passed as input streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID(), STREAM_ID), "true", - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, STREAM_ID), "true", + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertFalse(streamConfig.getDeleteCommittedMessages(PHYSICAL_STREAM)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID(), STREAM_ID), "true"))); + String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, STREAM_ID), "true"))); assertTrue(streamConfig.getDeleteCommittedMessages(STREAM_ID)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID(), STREAM_ID), "false"))); + String.format(StreamConfig.DELETE_COMMITTED_MESSAGES_FOR_STREAM_ID, STREAM_ID), "false"))); assertFalse(streamConfig.getDeleteCommittedMessages(STREAM_ID)); } @@ -459,48 +449,45 @@ public void testGetIsBounded() { // ignore mapping for other stream ids StreamConfig streamConfig = new StreamConfig(new MapConfig( - ImmutableMap.of(String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID(), OTHER_STREAM_ID), + ImmutableMap.of(String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID, OTHER_STREAM_ID), "true"))); assertFalse(streamConfig.getIsBounded(STREAM_ID)); // do not use stream id property if physical name is passed as input streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID(), STREAM_ID), "true", - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))); + String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID, STREAM_ID), "true", + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))); assertFalse(streamConfig.getIsBounded(PHYSICAL_STREAM)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID(), STREAM_ID), "true"))); + String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID, STREAM_ID), "true"))); assertTrue(streamConfig.getIsBounded(STREAM_ID)); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID(), STREAM_ID), "false"))); + String.format(StreamConfig.IS_BOUNDED_FOR_STREAM_ID, STREAM_ID), "false"))); assertFalse(streamConfig.getIsBounded(STREAM_ID)); } @Test public void testGetStreamIds() { assertEquals(ImmutableList.of(), ImmutableList.copyOf( - JavaConverters.asJavaIterableConverter(new StreamConfig(new MapConfig()).getStreamIds()).asJava())); + new StreamConfig(new MapConfig()).getStreamIds())); StreamConfig streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property", "value"))); - assertEquals(ImmutableList.of(STREAM_ID), - ImmutableList.copyOf(JavaConverters.asJavaIterableConverter(streamConfig.getStreamIds()).asJava())); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property", "value"))); + assertEquals(ImmutableSet.of(STREAM_ID), ImmutableSet.copyOf(streamConfig.getStreamIds())); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property.subProperty", "value"))); - assertEquals(ImmutableList.of(STREAM_ID), - ImmutableList.copyOf(JavaConverters.asJavaIterableConverter(streamConfig.getStreamIds()).asJava())); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property.subProperty", "value"))); + assertEquals(ImmutableSet.of(STREAM_ID), ImmutableSet.copyOf(streamConfig.getStreamIds())); streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property0", "value", - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property1", "value", - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property.subProperty0", "value", - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + ".property.subProperty1", "value", - String.format(StreamConfig.STREAM_ID_PREFIX(), OTHER_STREAM_ID) + ".property", "value"))); - assertEquals(ImmutableList.of(STREAM_ID, OTHER_STREAM_ID), - ImmutableList.copyOf(JavaConverters.asJavaIterableConverter(streamConfig.getStreamIds()).asJava())); + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property0", "value", + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property1", "value", + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property.subProperty0", "value", + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + ".property.subProperty1", "value", + String.format(StreamConfig.STREAM_ID_PREFIX, OTHER_STREAM_ID) + ".property", "value"))); + assertEquals(ImmutableSet.of(STREAM_ID, OTHER_STREAM_ID), ImmutableSet.copyOf(streamConfig.getStreamIds())); } private static void doTestSamzaProperty(String propertyName, String propertyValue, SamzaPropertyAssertion assertion) { @@ -518,21 +505,21 @@ private static void doTestSamzaProperty(String propertyName, String propertyValu private static void doTestSamzaPropertyAccess(String propertyName, String value, SamzaPropertyAssertion assertion) { // streams.. assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, value, // all streams need to have a system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM))), SYSTEM_STREAM); // systems..streams.. where stream id has no specified physical stream assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // specify the system for the stream id - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))), + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))), SYSTEM_STREAM); // systems..streams.. where stream is the streamId, system is from job default system assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // use job default system to get the system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM))), SYSTEM_STREAM); @@ -541,30 +528,30 @@ private static void doTestSamzaPropertyAccess(String propertyName, String value, assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, value, // specify the system for the stream id - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM))), + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM))), SYSTEM_STREAM); // systems..streams.. where no system mapping (fall back to SystemStream) assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); // systems..default.stream. where no system mapping (fall back to SystemStream) assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, value, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); // systems..streams.. where stream id has a physical name but property is from stream id assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // use job default system to get the system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); // systems..default.stream. where stream id has a physical name but property is from stream id @@ -573,7 +560,7 @@ private static void doTestSamzaPropertyAccess(String propertyName, String value, // use job default system to get the system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); } @@ -584,38 +571,38 @@ private static void doTestSamzaPropertyAccessWithPhysicalStream(String propertyN SamzaPropertyAssertion assertion) { // streams.. assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, value, // all streams need to have a system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM_PHYSICAL); // systems..streams.. with a specific system for the stream id assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, PHYSICAL_STREAM) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, PHYSICAL_STREAM) + propertyName, value, // specify the system for the stream id - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM_PHYSICAL); // systems..streams.. with system coming from job default system assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, PHYSICAL_STREAM) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, PHYSICAL_STREAM) + propertyName, value, // use job default system to get the system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM_PHYSICAL); // systems..default.stream. assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, value, // specify the system for the stream id - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM_PHYSICAL); } @@ -627,18 +614,18 @@ private static void doTestSamzaPropertyAccessWithPhysicalStream(String propertyN private static void doTestSamzaPropertyPriority(String propertyName, String value, SamzaPropertyAssertion assertion) { // streams.. vs. systems..streams.. assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, value, // all streams need to have a system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // this config should not be used - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE))), + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE))), SYSTEM_STREAM); // systems..streams.. vs. systems..default.stream. assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // specify the system for the stream id - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, // this config should not be used String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, UNUSED_VALUE))), SYSTEM_STREAM); @@ -654,7 +641,7 @@ private static void doTestSamzaPropertyPriority(String propertyName, String valu * systems..default.stream. without a system mapping in the config */ assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, value, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, value, // this config should not be used String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, UNUSED_VALUE))), SYSTEM_STREAM); @@ -669,16 +656,16 @@ private static void doTestSamzaPropertyMultipleStreams(String propertyName, Stri assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( String.format(SystemConfig.SYSTEM_DEFAULT_STREAMS_PREFIX_FORMAT, SYSTEM) + propertyName, value, // specify the systems for the streams - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID), SYSTEM, - String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), OTHER_STREAM_ID), SYSTEM))), + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID), SYSTEM, + String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, OTHER_STREAM_ID), SYSTEM))), SYSTEM_STREAM); } private static void doTestSamzaPropertyInvalidConfig(SamzaPropertyLookup lookup) { // configure physical stream and have mapping from stream id to physical stream StreamConfig streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM, - String.format(StreamConfig.STREAM_ID_PREFIX(), PHYSICAL_STREAM) + ".property", "value", + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM, + String.format(StreamConfig.STREAM_ID_PREFIX, PHYSICAL_STREAM) + ".property", "value", JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM))); try { lookup.doLookup(streamConfig, SYSTEM_STREAM_PHYSICAL); @@ -689,8 +676,8 @@ private static void doTestSamzaPropertyInvalidConfig(SamzaPropertyLookup lookup) // two separate stream ids map to same physical stream streamConfig = new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM, - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), OTHER_STREAM_ID), PHYSICAL_STREAM, + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM, + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, OTHER_STREAM_ID), PHYSICAL_STREAM, JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM))); try { lookup.doLookup(streamConfig, SYSTEM_STREAM_PHYSICAL); @@ -703,18 +690,18 @@ private static void doTestSamzaPropertyInvalidConfig(SamzaPropertyLookup lookup) private static void doTestSamzaPropertyDoesNotExist(String propertyName, SamzaPropertyAssertion assertion) { assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( // just put in some value which will be ignored - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))), + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + SAMZA_IGNORED_PROPERTY, UNUSED_VALUE))), SYSTEM_STREAM); /* * Won't use streams.. if streamId is mapped to a physical stream */ assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_ID_PREFIX(), STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.STREAM_ID_PREFIX, STREAM_ID) + propertyName, UNUSED_VALUE, // all streams need to have a system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); /* @@ -722,11 +709,11 @@ private static void doTestSamzaPropertyDoesNotExist(String propertyName, SamzaPr * is only specified using the physical stream */ assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, PHYSICAL_STREAM) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, PHYSICAL_STREAM) + propertyName, UNUSED_VALUE, // all streams need to have a system JobConfig.JOB_DEFAULT_SYSTEM, SYSTEM, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM); /* @@ -734,9 +721,9 @@ private static void doTestSamzaPropertyDoesNotExist(String propertyName, SamzaPr * system mapping for the stream id */ assertion.doAssertion(new StreamConfig(new MapConfig(ImmutableMap.of( - String.format(StreamConfig.STREAM_PREFIX(), SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, + String.format(StreamConfig.STREAM_PREFIX, SYSTEM, STREAM_ID) + propertyName, UNUSED_VALUE, // map the stream id to the physical stream - String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID), PHYSICAL_STREAM))), + String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID), PHYSICAL_STREAM))), SYSTEM_STREAM_PHYSICAL); } diff --git a/samza-core/src/test/java/org/apache/samza/config/TestSystemConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestSystemConfig.java index 5512bbcaa8..f8888f6ae3 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestSystemConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestSystemConfig.java @@ -26,6 +26,7 @@ import org.apache.samza.system.SystemFactory; import org.apache.samza.system.SystemProducer; import org.junit.Test; + import java.util.HashMap; import java.util.Map; @@ -157,13 +158,13 @@ public void testGetSystemKeySerde() { // value specified explicitly Config config = - new MapConfig(ImmutableMap.of(defaultStreamPrefixSystem1 + StreamConfig.KEY_SERDE(), system1KeySerde)); + new MapConfig(ImmutableMap.of(defaultStreamPrefixSystem1 + StreamConfig.KEY_SERDE, system1KeySerde)); SystemConfig systemConfig = new SystemConfig(config); assertEquals(system1KeySerde, systemConfig.getSystemKeySerde(MOCK_SYSTEM_NAME1).get()); // default stream property is unspecified, try fall back config key config = new MapConfig( - ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE(), + ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE, system1KeySerde)); systemConfig = new SystemConfig(config); assertEquals(system1KeySerde, systemConfig.getSystemKeySerde(MOCK_SYSTEM_NAME1).get()); @@ -171,15 +172,15 @@ public void testGetSystemKeySerde() { // default stream property is empty string, try fall back config key config = new MapConfig(ImmutableMap.of( // default stream property is empty - defaultStreamPrefixSystem1 + StreamConfig.KEY_SERDE(), "", + defaultStreamPrefixSystem1 + StreamConfig.KEY_SERDE, "", // fall back entry - String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE(), system1KeySerde)); + String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE, system1KeySerde)); systemConfig = new SystemConfig(config); assertEquals(system1KeySerde, systemConfig.getSystemKeySerde(MOCK_SYSTEM_NAME1).get()); // default stream property is unspecified, fall back is also empty config = new MapConfig( - ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE(), + ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.KEY_SERDE, "")); systemConfig = new SystemConfig(config); assertFalse(systemConfig.getSystemKeySerde(MOCK_SYSTEM_NAME1).isPresent()); @@ -197,13 +198,13 @@ public void testGetSystemMsgSerde() { // value specified explicitly Config config = - new MapConfig(ImmutableMap.of(defaultStreamPrefixSystem1 + StreamConfig.MSG_SERDE(), system1MsgSerde)); + new MapConfig(ImmutableMap.of(defaultStreamPrefixSystem1 + StreamConfig.MSG_SERDE, system1MsgSerde)); SystemConfig systemConfig = new SystemConfig(config); assertEquals(system1MsgSerde, systemConfig.getSystemMsgSerde(MOCK_SYSTEM_NAME1).get()); // default stream property is unspecified, try fall back config msg config = new MapConfig( - ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE(), + ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE, system1MsgSerde)); systemConfig = new SystemConfig(config); assertEquals(system1MsgSerde, systemConfig.getSystemMsgSerde(MOCK_SYSTEM_NAME1).get()); @@ -211,15 +212,15 @@ public void testGetSystemMsgSerde() { // default stream property is empty string, try fall back config msg config = new MapConfig(ImmutableMap.of( // default stream property is empty - defaultStreamPrefixSystem1 + StreamConfig.MSG_SERDE(), "", + defaultStreamPrefixSystem1 + StreamConfig.MSG_SERDE, "", // fall back entry - String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE(), system1MsgSerde)); + String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE, system1MsgSerde)); systemConfig = new SystemConfig(config); assertEquals(system1MsgSerde, systemConfig.getSystemMsgSerde(MOCK_SYSTEM_NAME1).get()); // default stream property is unspecified, fall back is also empty config = new MapConfig( - ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE(), + ImmutableMap.of(String.format(SystemConfig.SYSTEM_ID_PREFIX, MOCK_SYSTEM_NAME1) + StreamConfig.MSG_SERDE, "")); systemConfig = new SystemConfig(config); assertFalse(systemConfig.getSystemMsgSerde(MOCK_SYSTEM_NAME1).isPresent()); diff --git a/samza-core/src/test/java/org/apache/samza/execution/TestExecutionPlanner.java b/samza-core/src/test/java/org/apache/samza/execution/TestExecutionPlanner.java index 6a65286aa2..63d290a17a 100644 --- a/samza-core/src/test/java/org/apache/samza/execution/TestExecutionPlanner.java +++ b/samza-core/src/test/java/org/apache/samza/execution/TestExecutionPlanner.java @@ -40,7 +40,7 @@ import org.apache.samza.config.Config; import org.apache.samza.config.JobConfig; import org.apache.samza.config.MapConfig; -import org.apache.samza.config.StreamConfig$; +import org.apache.samza.config.StreamConfig; import org.apache.samza.config.TaskConfig; import org.apache.samza.serializers.StringSerde; import org.apache.samza.system.descriptors.GenericInputDescriptor; @@ -640,7 +640,7 @@ public void testDefaultPartitions() { @Test public void testBroadcastConfig() { Map map = new HashMap<>(config); - map.put(String.format(StreamConfig$.MODULE$.BROADCAST_FOR_STREAM_ID(), "input1"), "true"); + map.put(String.format(StreamConfig.BROADCAST_FOR_STREAM_ID, "input1"), "true"); Config cfg = new MapConfig(map); ExecutionPlanner planner = new ExecutionPlanner(cfg, streamManager); diff --git a/samza-core/src/test/java/org/apache/samza/testUtils/StreamTestUtils.java b/samza-core/src/test/java/org/apache/samza/testUtils/StreamTestUtils.java index 0aebd30aa2..6ad245c383 100644 --- a/samza-core/src/test/java/org/apache/samza/testUtils/StreamTestUtils.java +++ b/samza-core/src/test/java/org/apache/samza/testUtils/StreamTestUtils.java @@ -19,7 +19,7 @@ package org.apache.samza.testUtils; import java.util.Map; -import org.apache.samza.config.StreamConfig$; +import org.apache.samza.config.StreamConfig; public class StreamTestUtils { @@ -33,7 +33,7 @@ public class StreamTestUtils { */ public static void addStreamConfigs(Map configs, String streamId, String systemName, String physicalName) { - configs.put(String.format(StreamConfig$.MODULE$.SYSTEM_FOR_STREAM_ID(), streamId), systemName); - configs.put(String.format(StreamConfig$.MODULE$.PHYSICAL_NAME_FOR_STREAM_ID(), streamId), physicalName); + configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, streamId), systemName); + configs.put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, streamId), physicalName); } } \ No newline at end of file diff --git a/samza-core/src/test/java/org/apache/samza/util/TestStreamUtil.java b/samza-core/src/test/java/org/apache/samza/util/TestStreamUtil.java index 341eb24818..d7db63ab33 100644 --- a/samza-core/src/test/java/org/apache/samza/util/TestStreamUtil.java +++ b/samza-core/src/test/java/org/apache/samza/util/TestStreamUtil.java @@ -18,8 +18,6 @@ */ package org.apache.samza.util; -import java.util.HashMap; -import java.util.Map; import org.apache.samza.config.Config; import org.apache.samza.config.JobConfig; import org.apache.samza.config.MapConfig; @@ -27,6 +25,9 @@ import org.apache.samza.system.StreamSpec; import org.junit.Test; +import java.util.HashMap; +import java.util.Map; + import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNull; @@ -49,8 +50,8 @@ public class TestStreamUtil { @Test public void testGetStreamWithPhysicalNameInConfig() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -62,7 +63,7 @@ public void testGetStreamWithPhysicalNameInConfig() { @Test public void testGetStreamWithoutPhysicalNameInConfig() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -73,8 +74,8 @@ public void testGetStreamWithoutPhysicalNameInConfig() { @Test public void testGetStreamWithSystemAtStreamScopeInConfig() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -85,7 +86,7 @@ public void testGetStreamWithSystemAtStreamScopeInConfig() { @Test public void testGetStreamWithSystemAtDefaultScopeInConfig() { Config config = addConfigs(buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME), + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME), JobConfig.JOB_DEFAULT_SYSTEM, TEST_DEFAULT_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -97,8 +98,8 @@ public void testGetStreamWithSystemAtDefaultScopeInConfig() { @Test public void testGetStreamWithSystemAtBothScopesInConfig() { Config config = addConfigs(buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM), + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM), JobConfig.JOB_DEFAULT_SYSTEM, TEST_DEFAULT_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -110,7 +111,7 @@ public void testGetStreamWithSystemAtBothScopesInConfig() { @Test(expected = IllegalArgumentException.class) public void testGetStreamWithOutSystemInConfig() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -121,8 +122,8 @@ public void testGetStreamWithOutSystemInConfig() { @Test public void testGetStreamPropertiesPassthrough() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM, + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM, "systemProperty1", "systemValue1", "systemProperty2", "systemValue2", "systemProperty3", "systemValue3"); @@ -143,8 +144,8 @@ public void testGetStreamPropertiesPassthrough() { @Test public void testGetStreamSamzaPropertiesOmitted() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM, + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM, "systemProperty1", "systemValue1", "systemProperty2", "systemValue2", "systemProperty3", "systemValue3"); @@ -153,18 +154,18 @@ public void testGetStreamSamzaPropertiesOmitted() { Map properties = spec.getConfig(); assertEquals(3, properties.size()); - assertNull(properties.get(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID))); - assertNull(properties.get(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID))); - assertNull(spec.get(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), STREAM_ID))); - assertNull(spec.get(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), STREAM_ID))); + assertNull(properties.get(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID))); + assertNull(properties.get(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID))); + assertNull(spec.get(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, STREAM_ID))); + assertNull(spec.get(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, STREAM_ID))); } @Test public void testStreamConfigOverrides() { final String sysStreamPrefix = String.format("systems.%s.streams.%s.", TEST_SYSTEM, TEST_PHYSICAL_NAME); Config config = addConfigs(buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM, + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM, "systemProperty1", "systemValue1", "systemProperty2", "systemValue2", "systemProperty3", "systemValue3"), @@ -183,8 +184,8 @@ public void testStreamConfigOverrides() { @Test public void testStreamConfigOverridesWithSystemDefaults() { Config config = addConfigs(buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM, + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM, "segment.bytes", "5309"), String.format("systems.%s.default.stream.replication.factor", TEST_SYSTEM), "4", // System default property String.format("systems.%s.default.stream.segment.bytest", TEST_SYSTEM), "867" @@ -202,8 +203,8 @@ public void testStreamConfigOverridesWithSystemDefaults() { @Test public void testGetStreamPhysicalNameArgSimple() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME2, // This should be ignored because of the explicit arg - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME2, // This should be ignored because of the explicit arg + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -216,8 +217,8 @@ public void testGetStreamPhysicalNameArgSimple() { @Test public void testGetStreamPhysicalNameArgSpecialCharacters() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME_SPECIAL_CHARS, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME_SPECIAL_CHARS, + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); assertEquals(TEST_PHYSICAL_NAME_SPECIAL_CHARS, spec.getPhysicalName()); @@ -227,8 +228,8 @@ public void testGetStreamPhysicalNameArgSpecialCharacters() { @Test public void testGetStreamPhysicalNameArgNull() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), null, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.PHYSICAL_NAME, null, + StreamConfig.SYSTEM, TEST_SYSTEM); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); assertNull(spec.getPhysicalName()); @@ -238,8 +239,8 @@ public void testGetStreamPhysicalNameArgNull() { @Test public void testGetStreamSystemNameArgValid() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, // This should be ignored because of the explicit arg - StreamConfig.SYSTEM(), TEST_SYSTEM); // This too + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, // This should be ignored because of the explicit arg + StreamConfig.SYSTEM, TEST_SYSTEM); // This too StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); @@ -252,8 +253,8 @@ public void testGetStreamSystemNameArgValid() { @Test(expected = IllegalArgumentException.class) public void testGetStreamSystemNameArgInvalid() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), TEST_SYSTEM_INVALID); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, TEST_SYSTEM_INVALID); StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); } @@ -262,8 +263,8 @@ public void testGetStreamSystemNameArgInvalid() { @Test(expected = IllegalArgumentException.class) public void testGetStreamSystemNameArgEmpty() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), ""); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, ""); StreamSpec spec = StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); } @@ -272,8 +273,8 @@ public void testGetStreamSystemNameArgEmpty() { @Test(expected = IllegalArgumentException.class) public void testGetStreamSystemNameArgNull() { Config config = buildStreamConfig(STREAM_ID, - StreamConfig.PHYSICAL_NAME(), TEST_PHYSICAL_NAME, - StreamConfig.SYSTEM(), null); + StreamConfig.PHYSICAL_NAME, TEST_PHYSICAL_NAME, + StreamConfig.SYSTEM, null); StreamUtil.getStreamSpec(STREAM_ID, new StreamConfig(config)); } @@ -282,7 +283,7 @@ public void testGetStreamSystemNameArgNull() { @Test(expected = IllegalArgumentException.class) public void testGetStreamStreamIdInvalid() { Config config = buildStreamConfig(STREAM_ID_INVALID, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.SYSTEM, TEST_SYSTEM); StreamUtil.getStreamSpec(STREAM_ID_INVALID, new StreamConfig(config)); } @@ -291,7 +292,7 @@ public void testGetStreamStreamIdInvalid() { @Test(expected = IllegalArgumentException.class) public void testGetStreamStreamIdEmpty() { Config config = buildStreamConfig("", - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.SYSTEM, TEST_SYSTEM); StreamUtil.getStreamSpec("", new StreamConfig(config)); } @@ -300,7 +301,7 @@ public void testGetStreamStreamIdEmpty() { @Test(expected = IllegalArgumentException.class) public void testGetStreamStreamIdNull() { Config config = buildStreamConfig(null, - StreamConfig.SYSTEM(), TEST_SYSTEM); + StreamConfig.SYSTEM, TEST_SYSTEM); StreamUtil.getStreamSpec(null, new StreamConfig(config)); } @@ -311,7 +312,7 @@ public void testGetStreamStreamIdNull() { private Config buildStreamConfig(String streamId, String... kvs) { // inject streams.x. into each key for (int i = 0; i < kvs.length - 1; i += 2) { - kvs[i] = String.format(StreamConfig.STREAM_ID_PREFIX(), streamId) + kvs[i]; + kvs[i] = String.format(StreamConfig.STREAM_ID_PREFIX, streamId) + kvs[i]; } return buildConfig(kvs); } diff --git a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java index 1986bea54f..f0bce193d9 100644 --- a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java +++ b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java @@ -23,18 +23,6 @@ import com.google.common.base.Preconditions; import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Optional; -import java.util.Properties; -import java.util.Set; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.function.Function; -import java.util.stream.Collectors; import org.apache.commons.lang3.NotImplementedException; import org.apache.commons.lang3.StringUtils; import org.apache.kafka.clients.admin.AdminClient; @@ -78,12 +66,24 @@ import scala.Function0; import scala.Function1; import scala.Function2; -import scala.collection.JavaConverters; import scala.runtime.AbstractFunction0; import scala.runtime.AbstractFunction1; import scala.runtime.AbstractFunction2; import scala.runtime.BoxedUnit; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Function; +import java.util.stream.Collectors; + public class KafkaSystemAdmin implements SystemAdmin { private static final Logger LOG = LoggerFactory.getLogger(KafkaSystemAdmin.class); @@ -140,12 +140,13 @@ public KafkaSystemAdmin(String systemName, Config config, Consumer metadataConsu LOG.info("New admin client with props:" + props); adminClient = AdminClient.create(props); + StreamConfig streamConfig = new StreamConfig(config); + KafkaConfig kafkaConfig = new KafkaConfig(config); coordinatorStreamReplicationFactor = Integer.valueOf(kafkaConfig.getCoordinatorReplicationFactor()); coordinatorStreamProperties = getCoordinatorStreamProperties(kafkaConfig); - Map storeToChangelog = - JavaConverters.mapAsJavaMapConverter(kafkaConfig.getKafkaChangelogEnabledStores()).asJava(); + Map storeToChangelog = kafkaConfig.getKafkaChangelogEnabledStores(); // Construct the meta information for each topic, if the replication factor is not defined, // we use 2 (DEFAULT_REPL_FACTOR) as the number of replicas for the change log stream. changelogTopicMetaInformation = new HashMap<>(); @@ -688,8 +689,7 @@ static Map getIntermediateStreamProperties(Config config) { if (appConfig.getAppMode() == ApplicationConfig.ApplicationMode.BATCH) { StreamConfig streamConfig = new StreamConfig(config); - intermedidateStreamProperties = JavaConverters.asJavaCollectionConverter(streamConfig.getStreamIds()) - .asJavaCollection() + intermedidateStreamProperties = streamConfig.getStreamIds() .stream() .filter(streamConfig::getIsIntermediateStream) .collect(Collectors.toMap(Function.identity(), streamId -> { diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 993f0e4225..f8051f27f1 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -290,9 +290,10 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { } // The method returns a map of storenames to changelog topic names, which are configured to use kafka as the changelog stream - def getKafkaChangelogEnabledStores() = { + def getKafkaChangelogEnabledStores(): util.HashMap[String, String] = { val changelogConfigs = config.regexSubset(KafkaConfig.CHANGELOG_STREAM_NAMES_REGEX).asScala - var storeToChangelog = Map[String, String]() + //var storeToChangelog = Map[String, String]() + var storeToChangelog :util.HashMap[String, String] = new util.HashMap(); val storageConfig = new StorageConfig(config) val pattern = Pattern.compile(KafkaConfig.CHANGELOG_STREAM_NAMES_REGEX) @@ -304,7 +305,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { JavaOptionals.toRichOptional(storageConfig.getChangelogStream(storeName)).toOption.foreach(changelogName => { val systemStream = StreamUtil.getSystemStreamFromNames(changelogName) - storeToChangelog += storeName -> systemStream.getStream + storeToChangelog.put(storeName, systemStream.getStream) }) } storeToChangelog diff --git a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java index 88f8510a7f..aedf3f1607 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java +++ b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java @@ -99,6 +99,8 @@ public KafkaConsumerProxy(Consumer kafkaConsumer, String systemName, Strin /** * Add new partition to the list of polled partitions. * Bust only be called before {@link KafkaConsumerProxy#start} is called.. + * @param ssp - SystemStreamPartition to add + * @param nextOffset - add partition with this starting offset */ public void addTopicPartition(SystemStreamPartition ssp, long nextOffset) { LOG.info(String.format("Adding new topicPartition %s with offset %s to queue for consumer %s", ssp, nextOffset, @@ -339,6 +341,8 @@ protected IncomingMessageEnvelope handleNewRecord(ConsumerRecord consumerR /** * Protected to help extensions of this class build {@link IncomingMessageEnvelope}s. + * @param r consumer record to size + * @return the size of the serialized record */ protected int getRecordSize(ConsumerRecord r) { int keySize = (r.key() == null) ? 0 : r.serializedKeySize(); diff --git a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaSystemFactory.scala b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaSystemFactory.scala index deb80569ec..8d1fd6b8d3 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaSystemFactory.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaSystemFactory.scala @@ -28,6 +28,7 @@ import org.apache.samza.config.KafkaConfig.Config2Kafka import org.apache.samza.config._ import org.apache.samza.metrics.MetricsRegistry import org.apache.samza.system.{SystemAdmin, SystemConsumer, SystemFactory, SystemProducer} +import scala.collection.JavaConverters._ import org.apache.samza.util._ object KafkaSystemFactory extends Logging { @@ -109,7 +110,7 @@ class KafkaSystemFactory extends SystemFactory with Logging { val appConfig = new ApplicationConfig(config) if (appConfig.getAppMode == ApplicationMode.BATCH) { val streamConfig = new StreamConfig(config) - streamConfig.getStreamIds().filter(streamConfig.getIsIntermediateStream(_)).map(streamId => { + streamConfig.getStreamIds().asScala.filter(streamConfig.getIsIntermediateStream(_)).map(streamId => { // only the override here val properties = new Properties() properties.putIfAbsent("retention.ms", String.valueOf(KafkaConfig.DEFAULT_RETENTION_MS_FOR_BATCH)) diff --git a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala index d94f414b32..ea6c3f81e7 100644 --- a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala +++ b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala @@ -102,17 +102,17 @@ class TestKafkaConfig { assertEquals(kafkaConfig.getChangelogKafkaProperties("test2").getProperty("max.message.bytes"), "1024000") assertEquals(kafkaConfig.getChangelogKafkaProperties("test3").getProperty("cleanup.policy"), "compact") val storeToChangelog = kafkaConfig.getKafkaChangelogEnabledStores() - assertEquals("mychangelog1", storeToChangelog.get("test1").getOrElse("")) - assertEquals("mychangelog2", storeToChangelog.get("test2").getOrElse("")) - assertEquals("otherstream", storeToChangelog.get("test3").getOrElse("")) + assertEquals("mychangelog1", storeToChangelog.getOrDefault("test1", "")) + assertEquals("mychangelog2", storeToChangelog.getOrDefault("test2", "")) + assertEquals("otherstream", storeToChangelog.getOrDefault("test3", "")) assertNull(kafkaConfig.getChangelogKafkaProperties("test1").getProperty("retention.ms")) assertNull(kafkaConfig.getChangelogKafkaProperties("test2").getProperty("retention.ms")) props.setProperty("systems." + SYSTEM_NAME + ".samza.factory", "org.apache.samza.system.kafka.SomeOtherFactory") val storeToChangelog1 = kafkaConfig.getKafkaChangelogEnabledStores() - assertEquals("mychangelog1", storeToChangelog1.get("test1").getOrElse("")) - assertEquals("mychangelog2", storeToChangelog1.get("test2").getOrElse("")) - assertEquals("otherstream", storeToChangelog1.get("test3").getOrElse("")) + assertEquals("mychangelog1", storeToChangelog1.getOrDefault("test1", "")) + assertEquals("mychangelog2", storeToChangelog1.getOrDefault("test2", "")) + assertEquals("otherstream", storeToChangelog1.getOrDefault("test3", "")) assertEquals(kafkaConfig.getChangelogKafkaProperties("test4").getProperty("cleanup.policy"), "delete") assertEquals(kafkaConfig.getChangelogKafkaProperties("test4").getProperty("retention.ms"), "3600") diff --git a/samza-log4j/src/main/java/org/apache/samza/config/Log4jSystemConfig.java b/samza-log4j/src/main/java/org/apache/samza/config/Log4jSystemConfig.java index b449f90801..185d83b6e2 100644 --- a/samza-log4j/src/main/java/org/apache/samza/config/Log4jSystemConfig.java +++ b/samza-log4j/src/main/java/org/apache/samza/config/Log4jSystemConfig.java @@ -21,6 +21,8 @@ import org.apache.samza.system.SystemStream; +import java.util.Optional; + /** * This class contains the methods for getting properties that are needed by the @@ -82,7 +84,7 @@ public String getSerdeClass(String name) { public String getStreamSerdeName(String systemName, String streamName) { StreamConfig streamConfig = new StreamConfig(this); - scala.Option option = streamConfig.getStreamMsgSerde(new SystemStream(systemName, streamName)); - return option.isEmpty() ? null : option.get(); + Optional option = streamConfig.getStreamMsgSerde(new SystemStream(systemName, streamName)); + return option.isPresent() ? option.get() : null; } } diff --git a/samza-log4j2/src/main/java/org/apache/samza/config/Log4jSystemConfig.java b/samza-log4j2/src/main/java/org/apache/samza/config/Log4jSystemConfig.java index b449f90801..185d83b6e2 100644 --- a/samza-log4j2/src/main/java/org/apache/samza/config/Log4jSystemConfig.java +++ b/samza-log4j2/src/main/java/org/apache/samza/config/Log4jSystemConfig.java @@ -21,6 +21,8 @@ import org.apache.samza.system.SystemStream; +import java.util.Optional; + /** * This class contains the methods for getting properties that are needed by the @@ -82,7 +84,7 @@ public String getSerdeClass(String name) { public String getStreamSerdeName(String systemName, String streamName) { StreamConfig streamConfig = new StreamConfig(this); - scala.Option option = streamConfig.getStreamMsgSerde(new SystemStream(systemName, streamName)); - return option.isEmpty() ? null : option.get(); + Optional option = streamConfig.getStreamMsgSerde(new SystemStream(systemName, streamName)); + return option.isPresent() ? option.get() : null; } } diff --git a/samza-sql/src/main/java/org/apache/samza/sql/SamzaSqlInputMessage.java b/samza-sql/src/main/java/org/apache/samza/sql/SamzaSqlInputMessage.java index ad6706cb06..7ce03397cd 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/SamzaSqlInputMessage.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/SamzaSqlInputMessage.java @@ -43,9 +43,9 @@ private SamzaSqlInputMessage(KV keyAndMessageKV, SamzaSqlRelMsgM /** * Constructs a new SamzaSqMessage given the input arguments - * @param keyAndMessageKV - * @param metadata - * @return + * @param keyAndMessageKV key-value of the message + * @param metadata metadata of the message + * @return new object of SamzaSqlInputMessage type */ public static SamzaSqlInputMessage of (KV keyAndMessageKV, SamzaSqlRelMsgMetadata metadata) { return new SamzaSqlInputMessage(keyAndMessageKV, metadata); diff --git a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java index d92faaef96..0761f47f7b 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/interfaces/SqlIOConfig.java @@ -100,11 +100,11 @@ public SqlIOConfig(String systemName, String streamName, List sourcePart if (!isRemoteTable()) { // The below config is required for local table and streams but not for remote table. - streamConfigs.put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID(), streamId), streamName); + streamConfigs.put(String.format(StreamConfig.PHYSICAL_NAME_FOR_STREAM_ID, streamId), streamName); if (tableDescriptor != null) { // For local table, set the bootstrap config and default offset to oldest - streamConfigs.put(String.format(StreamConfig.BOOTSTRAP_FOR_STREAM_ID(), streamId), "true"); - streamConfigs.put(String.format(StreamConfig.CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID(), streamId), "oldest"); + streamConfigs.put(String.format(StreamConfig.BOOTSTRAP_FOR_STREAM_ID, streamId), "true"); + streamConfigs.put(String.format(StreamConfig.CONSUMER_OFFSET_DEFAULT_FOR_STREAM_ID, streamId), "oldest"); } } diff --git a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java index 9482a75d6c..97b5de93c4 100644 --- a/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java +++ b/samza-sql/src/main/java/org/apache/samza/sql/planner/SamzaSqlValidator.java @@ -63,7 +63,7 @@ public SamzaSqlValidator(Config config) { /** * Validate a list of sql statements * @param sqlStmts list of sql statements - * @throws SamzaSqlValidatorException + * @throws SamzaSqlValidatorException exception for sql validation */ public void validate(List sqlStmts) throws SamzaSqlValidatorException { SamzaSqlApplicationConfig sqlConfig = SamzaSqlDslConverter.getSqlConfig(sqlStmts, config); diff --git a/samza-test/src/main/java/org/apache/samza/test/framework/TestRunner.java b/samza-test/src/main/java/org/apache/samza/test/framework/TestRunner.java index 082b7278c2..afdbde7951 100644 --- a/samza-test/src/main/java/org/apache/samza/test/framework/TestRunner.java +++ b/samza-test/src/main/java/org/apache/samza/test/framework/TestRunner.java @@ -431,9 +431,9 @@ private void deleteDirectory(String path) { * over {@link org.apache.samza.application.descriptors.ApplicationDescriptor} generated configs */ private void addSerdeConfigs(StreamDescriptor descriptor) { - String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX(), descriptor.getStreamId()); - String keySerdeConfigKey = streamIdPrefix + StreamConfig.KEY_SERDE(); - String msgSerdeConfigKey = streamIdPrefix + StreamConfig.MSG_SERDE(); + String streamIdPrefix = String.format(StreamConfig.STREAM_ID_PREFIX, descriptor.getStreamId()); + String keySerdeConfigKey = streamIdPrefix + StreamConfig.KEY_SERDE; + String msgSerdeConfigKey = streamIdPrefix + StreamConfig.MSG_SERDE; this.configs.put(keySerdeConfigKey, null); this.configs.put(msgSerdeConfigKey, null); } diff --git a/samza-test/src/test/java/org/apache/samza/test/operator/TestAsyncFlatMap.java b/samza-test/src/test/java/org/apache/samza/test/operator/TestAsyncFlatMap.java index 4f5d510511..275be34945 100644 --- a/samza-test/src/test/java/org/apache/samza/test/operator/TestAsyncFlatMap.java +++ b/samza-test/src/test/java/org/apache/samza/test/operator/TestAsyncFlatMap.java @@ -105,7 +105,7 @@ public void testDownstreamOperatorExceptionIsBubbledUp() { } private List runTest(List pageViews, Map configs) { - configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), PAGE_VIEW_STREAM), TEST_SYSTEM); + configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, PAGE_VIEW_STREAM), TEST_SYSTEM); InMemorySystemDescriptor isd = new InMemorySystemDescriptor(TEST_SYSTEM); InMemoryInputDescriptor pageViewStreamDesc = isd diff --git a/samza-test/src/test/java/org/apache/samza/test/table/TestLocalTableWithSideInputsEndToEnd.java b/samza-test/src/test/java/org/apache/samza/test/table/TestLocalTableWithSideInputsEndToEnd.java index 2490707742..eecc6b4dd0 100644 --- a/samza-test/src/test/java/org/apache/samza/test/table/TestLocalTableWithSideInputsEndToEnd.java +++ b/samza-test/src/test/java/org/apache/samza/test/table/TestLocalTableWithSideInputsEndToEnd.java @@ -82,9 +82,9 @@ public void testJoinWithDurableSideInputTable() { private void runTest(String systemName, StreamApplication app, List pageViews, List profiles) { Map configs = new HashMap<>(); - configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), PAGEVIEW_STREAM), systemName); - configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), PROFILE_STREAM), systemName); - configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID(), ENRICHED_PAGEVIEW_STREAM), systemName); + configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, PAGEVIEW_STREAM), systemName); + configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, PROFILE_STREAM), systemName); + configs.put(String.format(StreamConfig.SYSTEM_FOR_STREAM_ID, ENRICHED_PAGEVIEW_STREAM), systemName); InMemorySystemDescriptor isd = new InMemorySystemDescriptor(systemName);