dawidwys commented on a change in pull request #6205: [FLINK-9642]Reduce the 
count to deal with state during a CEP process
URL: https://github.com/apache/flink/pull/6205#discussion_r209233134
 
 

 ##########
 File path: 
flink-libraries/flink-cep/src/main/java/org/apache/flink/cep/nfa/sharedbuffer/SharedBufferAccessor.java
 ##########
 @@ -0,0 +1,362 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOVICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  Vhe 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.flink.cep.nfa.sharedbuffer;
+
+import org.apache.flink.api.java.tuple.Tuple2;
+import org.apache.flink.cep.nfa.DeweyNumber;
+import org.apache.flink.util.WrappingRuntimeException;
+
+import org.apache.commons.lang3.StringUtils;
+
+import javax.annotation.Nullable;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Stack;
+
+import static 
org.apache.flink.cep.nfa.compiler.NFAStateNameHandler.getOriginalNameFromInternal;
+import static org.apache.flink.util.Preconditions.checkState;
+
+/**
+ * A shared buffer implementation which stores values under according state. 
Additionally, the values can be
+ * versioned such that it is possible to retrieve their predecessor element in 
the buffer.
+ *
+ * <p>The idea of the implementation is to have a buffer for incoming events 
with unique ids assigned to them. This way
+ * we do not need to deserialize events during processing and we store only 
one copy of the event.
+ *
+ * <p>The entries in {@link SharedBufferAccessor} are {@link 
SharedBufferNode}. The shared buffer node allows to store
+ * relations between different entries. A dewey versioning scheme allows to 
discriminate between
+ * different relations (e.g. preceding element).
+ *
+ * <p>The implementation is strongly based on the paper "Efficient Pattern 
Matching over Event Streams".
+ *
+ * @param <V> Type of the values
+ * @see <a 
href="https://people.cs.umass.edu/~yanlei/publications/sase-sigmod08.pdf";>
+ * https://people.cs.umass.edu/~yanlei/publications/sase-sigmod08.pdf</a>
+ */
+public class SharedBufferAccessor<V> implements AutoCloseable{
+
+       /** The cache of sharedBuffer.*/
+       private SharedBuffer<V> sharedBuffer;
+
+       public SharedBufferAccessor(SharedBuffer<V> sharedBuffer) {
+               this.sharedBuffer = sharedBuffer;
+       }
+
+       public void setSharedBuffer(SharedBuffer<V> sharedBuffer) {
+               this.sharedBuffer = sharedBuffer;
+       }
+
+       /**
+        * Notifies shared buffer that there will be no events with timestamp 
&lt;&eq; the given value. It allows to clear
+        * internal counters for number of events seen so far per timestamp.
+        *
+        * @param timestamp watermark, no earlier events will arrive
+        * @throws Exception Thrown if the system cannot access the state.
+        */
+       public void advanceTime(long timestamp) throws Exception {
+               sharedBuffer.advanceTime(timestamp);
+       }
+
+       /**
+        * Adds another unique event to the shared buffer and assigns a unique 
id for it. It automatically creates a
+        * lock on this event, so it won't be removed during processing of that 
event. Therefore the lock should be removed
+        * after processing all {@link 
org.apache.flink.cep.nfa.ComputationState}s
+        *
+        * <p><b>NOTE:</b>Should be called only once for each unique event!
+        *
+        * @param value event to be registered
+        * @return unique id of that event that should be used when putting 
entries to the buffer.
+        * @throws Exception Thrown if the system cannot access the state.
+        */
+       public EventId registerEvent(V value, long timestamp) throws Exception {
+               return sharedBuffer.registerEvent(value, timestamp);
+       }
+
+       /**
+        * Stores given value (value + timestamp) under the given state. It 
assigns a preceding element
+        * relation to the previous entry.
+        *
+        * @param stateName      name of the state that the event should be 
assigned to
+        * @param eventId        unique id of event assigned by this 
SharedBuffer
+        * @param previousNodeId id of previous entry (might be null if start 
of new run)
+        * @param version        Version of the previous relation
+        * @return assigned id of this element
+        * @throws Exception Thrown if the system cannot access the state.
+        */
+       public NodeId put(
+               final String stateName,
+               final EventId eventId,
+               @Nullable final NodeId previousNodeId,
+               final DeweyNumber version) throws Exception {
+
+               if (previousNodeId != null) {
+                       lockNode(previousNodeId);
+               }
+
+               NodeId currentNodeId = new NodeId(eventId, 
getOriginalNameFromInternal(stateName));
+               Lockable<SharedBufferNode> currentNode = 
sharedBuffer.getEntry(currentNodeId);
+               if (currentNode == null) {
+                       currentNode = new Lockable<>(new SharedBufferNode(), 0);
+                       lockEvent(eventId);
+               }
+
+               currentNode.getElement().addEdge(new SharedBufferEdge(
+                       previousNodeId,
+                       version));
+               sharedBuffer.cacheEntry(currentNodeId, currentNode);
 
 Review comment:
   Let's rename `cacheEntry -> upsertEntry`.

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
us...@infra.apache.org


With regards,
Apache Git Services

Reply via email to