This is an automated email from the ASF dual-hosted git repository.

sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git


The following commit(s) were added to refs/heads/master by this push:
     new 1318c86  Add key/value operations in StateContext (#1587)
1318c86 is described below

commit 1318c863c7efa7fbb0e3745a57ac06a52a3a8101
Author: Sijie Guo <[email protected]>
AuthorDate: Wed Jun 20 18:26:59 2018 -0700

    Add key/value operations in StateContext (#1587)
    
    * Add key/value operations in StateContext
    
    - add `put` operation to allow state to mutate key/value state
    - add `get` operation to allow fetching key/value state
    
    Signed-off-by: Sijie Guo <[email protected]>
---
 .../org/apache/pulsar/admin/cli/CmdFunctions.java  |   8 +-
 .../java/org/apache/pulsar/admin/cli/CmdSinks.java |   2 +-
 .../org/apache/pulsar/admin/cli/CmdSources.java    |   2 +-
 .../org/apache/pulsar/functions/api/Context.java   |  25 +++++
 .../pulsar/functions/instance/ContextImpl.java     |  45 +++++++-
 .../functions/instance/JavaInstanceRunnable.java   |  26 +----
 .../functions/instance/state/StateContext.java     |  44 +++++++-
 .../functions/instance/state/StateContextImpl.java |  40 ++++---
 .../pulsar/functions/instance/ContextImplTest.java | 122 +++++++++++++++++++++
 .../instance/state/StateContextImplTest.java       |  97 ++++++++++++++++
 .../example-stateful-function-config.yaml          |  30 +++++
 .../pulsar/functions/runtime/ProcessRuntime.java   |  11 +-
 .../functions/runtime/ProcessRuntimeFactory.java   |   7 +-
 .../functions/runtime/ProcessRuntimeTest.java      |   9 +-
 .../functions/worker/FunctionRuntimeManager.java   |   1 +
 15 files changed, 415 insertions(+), 54 deletions(-)

diff --git 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
index 1dbe50b..0850172 100644
--- 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
+++ 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
@@ -701,7 +701,7 @@ public class CmdFunctions extends CmdBase {
         @Override
         void runCmd() throws Exception {
             CmdFunctions.startLocalRun(convertProto2(functionConfig), 
functionConfig.getParallelism(),
-                    instanceIdOffset, brokerServiceUrl,
+                    instanceIdOffset, brokerServiceUrl, stateStorageServiceUrl,
                     
AuthenticationConfig.builder().clientAuthenticationPlugin(clientAuthPlugin)
                             
.clientAuthenticationParameters(clientAuthParams).useTls(useTls)
                             
.tlsAllowInsecureConnection(tlsAllowInsecureConnection)
@@ -791,7 +791,7 @@ public class CmdFunctions extends CmdBase {
             String tableNs = String.format(
                 "%s_%s",
                 tenant,
-                namespace);
+                namespace).replace('-', '_');
 
             String tableName = getFunctionName();
 
@@ -977,7 +977,7 @@ public class CmdFunctions extends CmdBase {
     }
 
     protected static void 
startLocalRun(org.apache.pulsar.functions.proto.Function.FunctionDetails 
functionDetails,
-            int parallelism, int instanceIdOffset, String brokerServiceUrl, 
AuthenticationConfig authConfig,
+            int parallelism, int instanceIdOffset, String brokerServiceUrl, 
String stateStorageServiceUrl, AuthenticationConfig authConfig,
             String userCodeFile, PulsarAdmin admin)
             throws Exception {
 
@@ -988,7 +988,7 @@ public class CmdFunctions extends CmdBase {
         if (serviceUrl == null) {
             serviceUrl = DEFAULT_SERVICE_URL;
         }
-        try (ProcessRuntimeFactory containerFactory = new 
ProcessRuntimeFactory(serviceUrl, authConfig, null, null,
+        try (ProcessRuntimeFactory containerFactory = new 
ProcessRuntimeFactory(serviceUrl, stateStorageServiceUrl, authConfig, null, 
null,
                 null)) {
             List<RuntimeSpawner> spawners = new LinkedList<>();
             for (int i = 0; i < parallelism; ++i) {
diff --git 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
index 83eb7f5..8d05ec8 100644
--- 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
+++ 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
@@ -129,7 +129,7 @@ public class CmdSinks extends CmdBase {
         @Override
         void runCmd() throws Exception {
             CmdFunctions.startLocalRun(createSinkConfigProto2(sinkConfig), 
sinkConfig.getParallelism(),
-                    0, brokerServiceUrl,
+                    0, brokerServiceUrl, null,
                     
AuthenticationConfig.builder().clientAuthenticationPlugin(clientAuthPlugin)
                             
.clientAuthenticationParameters(clientAuthParams).useTls(useTls)
                             
.tlsAllowInsecureConnection(tlsAllowInsecureConnection)
diff --git 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
index 32d9a6c..e49357c 100644
--- 
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
+++ 
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
@@ -128,7 +128,7 @@ public class CmdSources extends CmdBase {
         @Override
         void runCmd() throws Exception {
             CmdFunctions.startLocalRun(createSourceConfigProto2(sourceConfig), 
sourceConfig.getParallelism(),
-                    0, brokerServiceUrl,
+                    0, brokerServiceUrl, null,
                     
AuthenticationConfig.builder().clientAuthenticationPlugin(clientAuthPlugin)
                             
.clientAuthenticationParameters(clientAuthParams).useTls(useTls)
                             
.tlsAllowInsecureConnection(tlsAllowInsecureConnection)
diff --git 
a/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Context.java
 
b/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Context.java
index ff4afcf..2fa513e 100644
--- 
a/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Context.java
+++ 
b/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Context.java
@@ -18,6 +18,7 @@
  */
 package org.apache.pulsar.functions.api;
 
+import java.nio.ByteBuffer;
 import org.slf4j.Logger;
 
 import java.util.Collection;
@@ -114,6 +115,30 @@ public interface Context {
     void incrCounter(String key, long amount);
 
     /**
+     * Retrieve the counter value for the key.
+     *
+     * @param key name of the key
+     * @return the amount of the counter value for this key
+     */
+    long getCounter(String key);
+
+    /**
+     * Updare the state value for the key.
+     *
+     * @param key name of the key
+     * @param value state value of the key
+     */
+    void putState(String key, ByteBuffer value);
+
+    /**
+     * Retrieve the state value for the key.
+     *
+     * @param key name of the key
+     * @return the state value for the key.
+     */
+    ByteBuffer getState(String key);
+
+    /**
      * Get a map of all user-defined key/value configs for the function
      * @return The full map of user-defined config values
      */
diff --git 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
index b606402..fce440d 100644
--- 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
+++ 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
@@ -18,8 +18,11 @@
  */
 package org.apache.pulsar.functions.instance;
 
+import static com.google.common.base.Preconditions.checkState;
+
 import com.google.gson.Gson;
 import com.google.gson.reflect.TypeToken;
+import java.nio.ByteBuffer;
 import lombok.Getter;
 import lombok.Setter;
 import org.apache.commons.lang.StringUtils;
@@ -215,12 +218,48 @@ class ContextImpl implements Context {
         return userConfigs;
     }
 
+
+    private void ensureStateEnabled() {
+        checkState(null != stateContext, "State is not enabled.");
+    }
+
     @Override
     public void incrCounter(String key, long amount) {
-        if (null != stateContext) {
+        ensureStateEnabled();
+        try {
             stateContext.incr(key, amount);
-        } else {
-            throw new RuntimeException("State is not enabled.");
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to increment key '" + key + "' 
by amount '" + amount + "'", e);
+        }
+    }
+
+    @Override
+    public long getCounter(String key) {
+        ensureStateEnabled();
+        try {
+            return stateContext.getAmount(key);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to retrieve counter from key '" 
+ key + "'");
+        }
+    }
+
+    @Override
+    public void putState(String key, ByteBuffer value) {
+        ensureStateEnabled();
+        try {
+            stateContext.put(key, value);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to update the state value for 
key '" + key + "'");
+        }
+    }
+
+    @Override
+    public ByteBuffer getState(String key) {
+        ensureStateEnabled();
+        try {
+            return stateContext.getValue(key);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to retrieve the state value for 
key '" + key + "'");
         }
     }
 
diff --git 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
index c68910b..d538c8b 100644
--- 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
+++ 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
@@ -30,7 +30,6 @@ import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
-import java.util.concurrent.CompletableFuture;
 
 import lombok.AccessLevel;
 import lombok.Getter;
@@ -160,6 +159,10 @@ public class JavaInstanceRunnable implements 
AutoCloseable, Runnable {
     public void run() {
         try {
             javaInstance = setupJavaInstance();
+            if (null != stateTable) {
+                StateContextImpl stateContext = new 
StateContextImpl(stateTable);
+                javaInstance.getContext().setStateContext(stateContext);
+            }
             while (true) {
 
                 currentRecord = readInput();
@@ -171,16 +174,6 @@ public class JavaInstanceRunnable implements 
AutoCloseable, Runnable {
                     }
                 }
 
-                // state object is per function, because we need to have the 
ability to know what updates
-                // are made in this function and ensure we only acknowledge 
after the state is persisted.
-                StateContextImpl stateContext;
-                if (null != stateTable) {
-                    stateContext = new StateContextImpl(stateTable);
-                    javaInstance.getContext().setStateContext(stateContext);
-                } else {
-                    stateContext = null;
-                }
-
                 // process the message
                 long processAt = System.currentTimeMillis();
                 stats.incrementProcessed(processAt);
@@ -201,16 +194,6 @@ public class JavaInstanceRunnable implements 
AutoCloseable, Runnable {
                 long doneProcessing = System.currentTimeMillis();
                 log.debug("Got result: {}", result.getResult());
 
-                if (null != stateContext) {
-                    CompletableFuture completableFuture = stateContext.flush();
-
-                    try {
-                        completableFuture.join();
-                    } catch (Exception e) {
-                        log.error("Failed to flush the state updates of 
message {}", currentRecord, e);
-                        currentRecord.fail();
-                    }
-                }
                 try {
                     processResult(currentRecord, result, processAt, 
doneProcessing);
                 } catch (Exception e) {
@@ -259,7 +242,6 @@ public class JavaInstanceRunnable implements AutoCloseable, 
Runnable {
         ).replace('-', '_');
         String tableName = instanceConfig.getFunctionDetails().getName();
 
-        // TODO (sijie): use endpoint for now
         StorageClientSettings settings = StorageClientSettings.newBuilder()
                 .serviceUri(stateStorageServiceUrl)
                 .clientName("function-" + tableNs + "/" + tableName)
diff --git 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContext.java
 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContext.java
index 83470f5..c17e8b6 100644
--- 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContext.java
+++ 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContext.java
@@ -18,15 +18,53 @@
  */
 package org.apache.pulsar.functions.instance.state;
 
-import java.util.concurrent.CompletableFuture;
+import java.nio.ByteBuffer;
 
 /**
  * A state context per function.
  */
 public interface StateContext {
 
-    void incr(String key, long amount);
+    /**
+     * Increment the given <i>key</i> by the given <i>amount</i>.
+     *
+     * @param key key to increment
+     * @param amount the amount incremented
+     */
+    void incr(String key, long amount) throws Exception;
 
-    CompletableFuture<Void> flush();
+    /**
+     * Update the given <i>key</i> to the provide <i>value</i>.
+     *
+     * <p>NOTE: the put operation might or might not be applied directly to 
the global state until
+     * the state is flushed via {@link #flush()} at the completion of function 
execution.
+     *
+     * <p>The behavior of `PUT` is non-deterministic, if two function 
instances attempt to update
+     * same key around the same time, there is no guarantee which update will 
be the final result.
+     * That says, if you attempt to get amount via {@link #getAmount(String)}, 
increment the amount
+     * based on the function computation logic, and update the computed amount 
back. one update will
+     * overwrite the other update. For this case, you are encouraged to use 
{@link #incr(String, long)}
+     * instead.
+     *
+     * @param key key to update.
+     * @param value value to update
+     */
+    void put(String key, ByteBuffer value) throws Exception;
+
+    /**
+     * Get the value of a given <i>key</i>.
+     *
+     * @param key key to retrieve
+     * @return a completable future representing the retrieve result.
+     */
+    ByteBuffer getValue(String key) throws Exception;
+
+    /**
+     * Get the amount of a given <i>key</i>.
+     *
+     * @param key key to retrieve
+     * @return a completable future representing the retrieve result.
+     */
+    long getAmount(String key) throws Exception;
 
 }
diff --git 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContextImpl.java
 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContextImpl.java
index a5ff408..1a2c26d 100644
--- 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContextImpl.java
+++ 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/state/StateContextImpl.java
@@ -19,14 +19,12 @@
 package org.apache.pulsar.functions.instance.state;
 
 import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.apache.bookkeeper.common.concurrent.FutureUtils.result;
 
 import io.netty.buffer.ByteBuf;
 import io.netty.buffer.Unpooled;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.concurrent.CompletableFuture;
+import java.nio.ByteBuffer;
 import org.apache.bookkeeper.api.kv.Table;
-import org.apache.bookkeeper.common.concurrent.FutureUtils;
 
 /**
  * This class accumulates the state updates from one function.
@@ -36,27 +34,41 @@ import org.apache.bookkeeper.common.concurrent.FutureUtils;
 public class StateContextImpl implements StateContext {
 
     private final Table<ByteBuf, ByteBuf> table;
-    // the list
-    private final List<CompletableFuture<Void>> updates;
 
     public StateContextImpl(Table<ByteBuf, ByteBuf> table) {
         this.table = table;
-        this.updates = new ArrayList<>();
     }
 
     @Override
-    public void incr(String key, long amount) {
+    public void incr(String key, long amount) throws Exception {
         // TODO: this can be optimized with a batch operation.
-        updates.add(table.increment(
+        result(table.increment(
             Unpooled.wrappedBuffer(key.getBytes(UTF_8)),
             amount));
     }
 
-    /**
-     * flush and wait all the updates to be completed.
-     */
     @Override
-    public CompletableFuture<Void> flush() {
-        return FutureUtils.collect(updates).thenApply(ignored -> null);
+    public void put(String key, ByteBuffer value) throws Exception {
+        result(table.put(
+            Unpooled.wrappedBuffer(key.getBytes(UTF_8)),
+            Unpooled.wrappedBuffer(value)));
+    }
+
+    @Override
+    public ByteBuffer getValue(String key) throws Exception {
+        ByteBuf data = 
result(table.get(Unpooled.wrappedBuffer(key.getBytes(UTF_8))));
+        try {
+            ByteBuffer result = ByteBuffer.allocate(data.readableBytes());
+            data.readBytes(result);
+            return result;
+        } finally {
+            data.release();
+        }
     }
+
+    @Override
+    public long getAmount(String key) throws Exception {
+        return 
result(table.getNumber(Unpooled.wrappedBuffer(key.getBytes(UTF_8))));
+    }
+
 }
diff --git 
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/ContextImplTest.java
 
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/ContextImplTest.java
new file mode 100644
index 0000000..0fb027e
--- /dev/null
+++ 
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/ContextImplTest.java
@@ -0,0 +1,122 @@
+/**
+ * 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.pulsar.functions.instance;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.mockito.Matchers.eq;
+import static org.mockito.Matchers.same;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+import java.nio.ByteBuffer;
+import org.apache.pulsar.client.api.Consumer;
+import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.functions.instance.state.StateContextImpl;
+import org.apache.pulsar.functions.proto.Function.FunctionDetails;
+import org.junit.Before;
+import org.junit.Test;
+import org.slf4j.Logger;
+
+/**
+ * Unit test {@link ContextImpl}.
+ */
+public class ContextImplTest {
+
+    private InstanceConfig config;
+    private Logger logger;
+    private PulsarClient client;
+    private ClassLoader classLoader;
+    private Consumer consumer;
+    private ContextImpl context;
+
+    @Before
+    public void setup() {
+        config = new InstanceConfig();
+        FunctionDetails functionDetails = FunctionDetails.newBuilder()
+            .setUserConfig("")
+            .build();
+        config.setFunctionDetails(functionDetails);
+        logger = mock(Logger.class);
+        client = mock(PulsarClient.class);
+        classLoader = getClass().getClassLoader();
+        consumer = mock(Consumer.class);
+        context = new ContextImpl(
+            config,
+            logger,
+            client,
+            classLoader,
+            consumer
+        );
+    }
+
+    @Test(expected = IllegalStateException.class)
+    public void testIncrCounterStateDisabled() {
+        context.incrCounter("test-key", 10);
+    }
+
+    @Test(expected = IllegalStateException.class)
+    public void testGetCounterStateDisabled() {
+        context.getCounter("test-key");
+    }
+
+    @Test(expected = IllegalStateException.class)
+    public void testPutStateStateDisabled() {
+        context.putState("test-key", 
ByteBuffer.wrap("test-value".getBytes(UTF_8)));
+    }
+
+    @Test(expected = IllegalStateException.class)
+    public void testGetStateStateDisabled() {
+        context.getState("test-key");
+    }
+
+    @Test
+    public void testIncrCounterStateEnabled() throws Exception {
+        StateContextImpl stateContext = mock(StateContextImpl.class);
+        context.setStateContext(stateContext);
+        context.incrCounter("test-key", 10L);
+        verify(stateContext, times(1)).incr(eq("test-key"), eq(10L));
+    }
+
+    @Test
+    public void testGetCounterStateEnabled() throws Exception {
+        StateContextImpl stateContext = mock(StateContextImpl.class);
+        context.setStateContext(stateContext);
+        context.getCounter("test-key");
+        verify(stateContext, times(1)).getAmount(eq("test-key"));
+    }
+
+    @Test
+    public void testPutStateStateEnabled() throws Exception {
+        StateContextImpl stateContext = mock(StateContextImpl.class);
+        context.setStateContext(stateContext);
+        ByteBuffer buffer = ByteBuffer.wrap("test-value".getBytes(UTF_8));
+        context.putState("test-key", buffer);
+        verify(stateContext, times(1)).put(eq("test-key"), same(buffer));
+    }
+
+    @Test
+    public void testGetStateStateEnabled() throws Exception {
+        StateContextImpl stateContext = mock(StateContextImpl.class);
+        context.setStateContext(stateContext);
+        context.getState("test-key");
+        verify(stateContext, times(1)).getValue(eq("test-key"));
+    }
+
+}
diff --git 
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/state/StateContextImplTest.java
 
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/state/StateContextImplTest.java
new file mode 100644
index 0000000..afe0403
--- /dev/null
+++ 
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/state/StateContextImplTest.java
@@ -0,0 +1,97 @@
+/**
+ * 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.pulsar.functions.instance.state;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.mockito.Matchers.any;
+import static org.mockito.Matchers.anyLong;
+import static org.mockito.Matchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+import static org.testng.AssertJUnit.assertEquals;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+import java.nio.ByteBuffer;
+import org.apache.bookkeeper.api.kv.Table;
+import org.apache.bookkeeper.common.concurrent.FutureUtils;
+import org.junit.Before;
+import org.junit.Test;
+
+/**
+ * Unit test {@link StateContextImpl}.
+ */
+public class StateContextImplTest {
+
+    private Table<ByteBuf, ByteBuf> mockTable;
+    private StateContextImpl stateContext;
+
+    @Before
+    public void setup() {
+        this.mockTable = mock(Table.class);
+        this.stateContext = new StateContextImpl(mockTable);
+    }
+
+    @Test
+    public void testIncr() throws Exception {
+        when(mockTable.increment(any(ByteBuf.class), anyLong()))
+            .thenReturn(FutureUtils.Void());
+        stateContext.incr("test-key", 10L);
+        verify(mockTable, times(1)).increment(
+            eq(Unpooled.copiedBuffer("test-key", UTF_8)),
+            eq(10L)
+        );
+    }
+
+    @Test
+    public void testPut() throws Exception {
+        when(mockTable.put(any(ByteBuf.class), any(ByteBuf.class)))
+            .thenReturn(FutureUtils.Void());
+        stateContext.put("test-key", 
ByteBuffer.wrap("test-value".getBytes(UTF_8)));
+        verify(mockTable, times(1)).put(
+            eq(Unpooled.copiedBuffer("test-key", UTF_8)),
+            eq(Unpooled.copiedBuffer("test-value", UTF_8))
+        );
+    }
+
+    @Test
+    public void testGetValue() throws Exception {
+        ByteBuf returnedValue = Unpooled.copiedBuffer("test-value", UTF_8);
+        when(mockTable.get(any(ByteBuf.class)))
+            .thenReturn(FutureUtils.value(returnedValue));
+        ByteBuffer result = stateContext.getValue("test-key");
+        assertEquals("test-value", new String(result.array(), UTF_8));
+        verify(mockTable, times(1)).get(
+            eq(Unpooled.copiedBuffer("test-key", UTF_8))
+        );
+    }
+
+    @Test
+    public void testGetAmount() throws Exception {
+        when(mockTable.getNumber(any(ByteBuf.class)))
+            .thenReturn(FutureUtils.value(10L));
+        assertEquals(10L, stateContext.getAmount("test-key"));
+        verify(mockTable, times(1)).getNumber(
+            eq(Unpooled.copiedBuffer("test-key", UTF_8))
+        );
+    }
+
+}
diff --git 
a/pulsar-functions/java-examples/src/main/resources/example-stateful-function-config.yaml
 
b/pulsar-functions/java-examples/src/main/resources/example-stateful-function-config.yaml
new file mode 100644
index 0000000..4c758c9
--- /dev/null
+++ 
b/pulsar-functions/java-examples/src/main/resources/example-stateful-function-config.yaml
@@ -0,0 +1,30 @@
+#
+# 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.
+#
+
+tenant: "test"
+namespace: "test-namespace"
+name: "stateful-example"
+className: "org.apache.pulsar.functions.api.examples.CounterFunction"
+inputs: ["test_stateful_src"]
+userConfig:
+  "PublishTopic": "test_stateful_result"
+
+output: "test_stateful_result"
+autoAck: true
+parallelism: 1
diff --git 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntime.java
 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntime.java
index e94e39d..736558b 100644
--- 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntime.java
+++ 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntime.java
@@ -70,10 +70,11 @@ class ProcessRuntime implements Runtime {
                    String logDirectory,
                    String codeFile,
                    String pulsarServiceUrl,
+                   String stateStorageServiceUrl,
                    AuthenticationConfig authConfig) {
         this.instanceConfig = instanceConfig;
         this.instancePort = instanceConfig.getPort();
-        this.processArgs = composeArgs(instanceConfig, instanceFile, 
logDirectory, codeFile, pulsarServiceUrl,
+        this.processArgs = composeArgs(instanceConfig, instanceFile, 
logDirectory, codeFile, pulsarServiceUrl, stateStorageServiceUrl,
                 authConfig);
     }
 
@@ -82,6 +83,7 @@ class ProcessRuntime implements Runtime {
                                      String logDirectory,
                                      String codeFile,
                                      String pulsarServiceUrl,
+                                     String stateStorageServiceUrl,
                                      AuthenticationConfig authConfig) {
         List<String> args = new LinkedList<>();
         if (instanceConfig.getFunctionDetails().getRuntime() == 
Function.FunctionDetails.Runtime.JAVA) {
@@ -228,6 +230,13 @@ class ProcessRuntime implements Runtime {
             args.add("--sink_serde_classname");
             
args.add(instanceConfig.getFunctionDetails().getSink().getSerDeClassName());
         }
+
+        // state storage configs
+        if (null != stateStorageServiceUrl
+            && instanceConfig.getFunctionDetails().getRuntime() == 
Function.FunctionDetails.Runtime.JAVA) {
+            args.add("--state_storage_serviceurl");
+            args.add(stateStorageServiceUrl);
+        }
         return args;
     }
 
diff --git 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntimeFactory.java
 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntimeFactory.java
index 8fc5b90..109d5ea 100644
--- 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntimeFactory.java
+++ 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/ProcessRuntimeFactory.java
@@ -33,7 +33,8 @@ import java.nio.file.Paths;
 @Slf4j
 public class ProcessRuntimeFactory implements RuntimeFactory {
 
-    private String pulsarServiceUrl;
+    private final String pulsarServiceUrl;
+    private final String stateStorageServiceUrl;
     private AuthenticationConfig authConfig;
     private String javaInstanceJarFile;
     private String pythonInstanceFile;
@@ -41,12 +42,13 @@ public class ProcessRuntimeFactory implements 
RuntimeFactory {
 
     @VisibleForTesting
     public ProcessRuntimeFactory(String pulsarServiceUrl,
+                                 String stateStorageServiceUrl,
                                  AuthenticationConfig authConfig,
                                  String javaInstanceJarFile,
                                  String pythonInstanceFile,
                                  String logDirectory) {
-
         this.pulsarServiceUrl = pulsarServiceUrl;
+        this.stateStorageServiceUrl = stateStorageServiceUrl;
         this.authConfig = authConfig;
         this.javaInstanceJarFile = javaInstanceJarFile;
         this.pythonInstanceFile = pythonInstanceFile;
@@ -106,6 +108,7 @@ public class ProcessRuntimeFactory implements 
RuntimeFactory {
             logDirectory,
             codeFile,
             pulsarServiceUrl,
+            stateStorageServiceUrl,
             authConfig);
     }
 
diff --git 
a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/ProcessRuntimeTest.java
 
b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/ProcessRuntimeTest.java
index 675f3de..9e40cb5 100644
--- 
a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/ProcessRuntimeTest.java
+++ 
b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/ProcessRuntimeTest.java
@@ -57,6 +57,7 @@ public class ProcessRuntimeTest {
     private final String javaInstanceJarFile;
     private final String pythonInstanceFile;
     private final String pulsarServiceUrl;
+    private final String stateStorageServiceUrl;
     private final String logDirectory;
 
     public ProcessRuntimeTest() {
@@ -64,9 +65,10 @@ public class ProcessRuntimeTest {
         this.javaInstanceJarFile = "/Users/user/JavaInstance.jar";
         this.pythonInstanceFile = "/Users/user/PythonInstance.py";
         this.pulsarServiceUrl = "pulsar://localhost:6670";
+        this.stateStorageServiceUrl = "bk://localhost:4181";
         this.logDirectory = "Users/user/logs";
         this.factory = new ProcessRuntimeFactory(
-            pulsarServiceUrl, null, javaInstanceJarFile, pythonInstanceFile, 
logDirectory);
+            pulsarServiceUrl, stateStorageServiceUrl, null, 
javaInstanceJarFile, pythonInstanceFile, logDirectory);
     }
 
     @AfterMethod
@@ -115,7 +117,7 @@ public class ProcessRuntimeTest {
 
         ProcessRuntime container = factory.createContainer(config, 
userJarFile);
         List<String> args = container.getProcessArgs();
-        assertEquals(args.size(), 53);
+        assertEquals(args.size(), 55);
         String expectedArgs = "java -cp " + javaInstanceJarFile + " 
-Dlog4j.configurationFile=java_instance_log4j2.yml "
                 + "-Dpulsar.log.dir=" + logDirectory + "/functions" + " 
-Dpulsar.log.file=" + config.getFunctionDetails().getName()
                 + " org.apache.pulsar.functions.runtime.JavaInstanceMain"
@@ -138,7 +140,8 @@ public class ProcessRuntimeTest {
                 + " --sink_classname " + 
config.getFunctionDetails().getSink().getClassName()
                 + " --sink_type_classname " + 
config.getFunctionDetails().getSink().getTypeClassName()
                 + " --sink_topic " + 
config.getFunctionDetails().getSink().getTopic()
-                + " --sink_serde_classname " + 
config.getFunctionDetails().getSink().getSerDeClassName();
+                + " --sink_serde_classname " + 
config.getFunctionDetails().getSink().getSerDeClassName()
+                + " --state_storage_serviceurl " + stateStorageServiceUrl;
         assertEquals(expectedArgs, String.join(" ", args));
     }
 
diff --git 
a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java
 
b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java
index 2512295..5c6184f 100644
--- 
a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java
+++ 
b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionRuntimeManager.java
@@ -109,6 +109,7 @@ public class FunctionRuntimeManager implements 
AutoCloseable{
         } else if (workerConfig.getProcessContainerFactory() != null) {
             this.runtimeFactory = new ProcessRuntimeFactory(
                     workerConfig.getPulsarServiceUrl(),
+                    workerConfig.getStateStorageServiceUrl(),
                     authConfig,
                     
workerConfig.getProcessContainerFactory().getJavaInstanceJarLocation(),
                     
workerConfig.getProcessContainerFactory().getPythonInstanceLocation(),

Reply via email to