d2990684a2bb6ab9aa3eac03f245aa22eeda11fd
[samza.git] / samza-api / src / main / java / org / apache / samza / operators / StreamGraph.java
1 /*
2 * Licensed to the Apache Software Foundation (ASF) under one
3 * or more contributor license agreements. See the NOTICE file
4 * distributed with this work for additional information
5 * regarding copyright ownership. The ASF licenses this file
6 * to you under the Apache License, Version 2.0 (the
7 * "License"); you may not use this file except in compliance
8 * with the License. You may obtain a copy of the License at
9 *
10 * http://www.apache.org/licenses/LICENSE-2.0
11 *
12 * Unless required by applicable law or agreed to in writing,
13 * software distributed under the License is distributed on an
14 * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15 * KIND, either express or implied. See the License for the
16 * specific language governing permissions and limitations
17 * under the License.
18 */
19 package org.apache.samza.operators;
20
21 import org.apache.samza.annotation.InterfaceStability;
22
23 import java.util.function.BiFunction;
24 import java.util.function.Function;
25
26 /**
27 * Provides APIs for accessing {@link MessageStream}s to be used to create the DAG of transforms.
28 */
29 @InterfaceStability.Unstable
30 public interface StreamGraph {
31
32 /**
33 * Gets the input {@link MessageStream} corresponding to the logical {@code streamId}. Multiple invocations of
34 * this method with the same {@code streamId} will throw an {@link IllegalStateException}
35 *
36 * @param streamId the unique logical ID for the stream
37 * @param msgBuilder the {@link BiFunction} to convert the incoming key and message to a message
38 * in the input {@link MessageStream}
39 * @param <K> the type of key in the incoming message
40 * @param <V> the type of message in the incoming message
41 * @param <M> the type of message in the input {@link MessageStream}
42 * @return the input {@link MessageStream}
43 * @throws IllegalStateException when invoked multiple times with the same {@code streamId}
44 */
45 <K, V, M> MessageStream<M> getInputStream(String streamId, BiFunction<? super K, ? super V, ? extends M> msgBuilder);
46
47 /**
48 * Gets the {@link OutputStream} corresponding to the logical {@code streamId}. Multiple invocations of
49 * this method with the same {@code streamId} will throw an {@link IllegalStateException}
50 *
51 * @param streamId the unique logical ID for the stream
52 * @param keyExtractor the {@link Function} to extract the outgoing key from the output message
53 * @param msgExtractor the {@link Function} to extract the outgoing message from the output message
54 * @param <K> the type of key in the outgoing message
55 * @param <V> the type of message in the outgoing message
56 * @param <M> the type of message in the {@link OutputStream}
57 * @return the output {@link MessageStream}
58 * @throws IllegalStateException when invoked multiple times with the same {@code streamId}
59 */
60 <K, V, M> OutputStream<K, V, M> getOutputStream(String streamId,
61 Function<? super M, ? extends K> keyExtractor, Function<? super M, ? extends V> msgExtractor);
62
63 /**
64 * Sets the {@link ContextManager} for this {@link StreamGraph}.
65 *
66 * The provided {@code contextManager} will be initialized before the transformation functions
67 * and can be used to setup shared context between them.
68 *
69 * @param contextManager the {@link ContextManager} to use for the {@link StreamGraph}
70 * @return the {@link StreamGraph} with the {@code contextManager} as its {@link ContextManager}
71 */
72 StreamGraph withContextManager(ContextManager contextManager);
73
74 }