[ 
https://issues.apache.org/jira/browse/NIP-47?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Mark Payne reassigned NIP-47:
-----------------------------

    Assignee: Mark Payne

> Allow Connectors to more effectively evolve their flows on startup
> ------------------------------------------------------------------
>
>                 Key: NIP-47
>                 URL: https://issues.apache.org/jira/browse/NIP-47
>             Project: NiFi Improvement Proposal
>          Issue Type: Improvement
>            Reporter: Mark Payne
>            Assignee: Mark Payne
>            Priority: Major
>
> h1. Motivation
> Connectors provide a way to programmatically control a full NiFi flow. This 
> is a very powerful concept. A missing component here, though, is the ability 
> to effectively manage the case where a startup of NiFi results in a new 
> version of the Connector, and the Connector has a new version of the flow 
> where Connections were removed.
> In such a case, the Connector is forced to maintain the old flow in some way 
> and cause the data to drain using the old version of the flow, and only then 
> migrate to the new version. This is awkward and requires significant effort 
> and is error-prone.
> h1. Proposal
> My proposal, then, is to offer the ability for a Connector to relocate the 
> FlowFiles that exist in a given (no-longer-existing) Connection to a new 
> Connection, possibly transforming the Attributes along the way.
> h2. API
> My intention would be to introduce a new method to the {{Connector}} 
> interface:
> {code:java}
> /**
>  * Provides the opportunity for a Connector to evolve its flow from an 
> earlier version to the current version
>  * by informing the framework how existing data should be handled.
>  *
>  * @param context the context that provides the ability to evolve the flow 
> from its pre-existing version to its current version
>  */
> default void evolveFlow(FlowEvolutionContext context) {
> }{code}
> Where {{FlowEvolutionContext}} is defined as:
> {code:java}
> /**
>  * The FlowEvolutionContext provides the mechanism by which Connectors are 
> able to evolve their flows from
>  * prior versions of the Connectors' flows to the current version.
>  */
> public interface FlowEvolutionContext {
>     /**
>      * Instructs the framework to rehome any FlowFiles that are currently 
> homed to the given removed Connection ID
>      * to the Connection with the given destination Connection ID. Note that 
> invoking this method does not itself necessarily
>      * transfer any FlowFiles between connections. Rather, it informs the 
> framework that the move needs to happen.
>      * It is up to the framework to make the determination when this should 
> happen. The framework may choose to move
>      * the FlowFiles inline during this method invocation, or it may choose 
> to do so later. The framework does guarantee,
>      * however, that this move happens before the Connector is able to start.
>      *
>      * @param removedConnectionId the ID of the Connection in the previous 
> flow
>      * @param destinationConnectionId the ID of the Connection in the current 
> flow
>      * @throws IllegalStateException if no Connection exists in the current 
> flow with the given ID,
>      *   or if a Connection still exists in the current flow with an ID 
> matching the removedConnectionId.
>      */
>     default void relocateFlowFiles(String removedConnectionId, String 
> destinationConnectionId) {
>         relocateFlowFiles(removedConnectionId, destinationConnectionId, 
> AttributesTransformer.NOP);
>     }
>     /**
>      * Instructs the framework to rehome any FlowFiles that are currently 
> homed to the given removed Connection ID
>      * to the Connection with the given destination Connection ID, 
> transforming the FlowFile's attributes as defined by the given
>      * transformer. Note that invoking this method does not itself necessarily
>      * transfer any FlowFiles between connections. Rather, it informs the 
> framework that the move needs to happen.
>      * It is up to the framework to make the determination when this should 
> happen. The framework may choose to move
>      * the FlowFiles inline during this method invocation, or it may choose 
> to do so later. The framework does guarantee,
>      * however, that this move happens before the Connector is able to start.
>      *
>      * @param removedConnectionId the ID of the Connection in the previous 
> flow
>      * @param destinationConnectionId the ID of the Connection in the current 
> flow
>      * @param attributesTransformer a transformer that can be used to 
> manipulate the attributes of FlowFiles that are moved
>      * @throws IllegalStateException if no Connection exists in the current 
> flow with the given ID,
>      *   or if a Connection still exists in the current flow with an ID 
> matching the removedConnectionId.
>      */
>     default void relocateFlowFiles(String removedConnectionId, String 
> destinationConnectionId, AttributesTransformer attributesTransformer) {
>         relocateFlowFiles(removedConnectionId, destinationConnectionId, 
> flowFile -> true, attributesTransformer);
>     }
>     /**
>      * Instructs the framework to rehome any FlowFiles that are currently 
> homed to the given removed Connection ID, and that match the given Predicate,
>      * to the Connection with the given destination Connection ID. The 
> FlowFiles' attributes will be transformed according to the provided 
> AttributesTransformer.
>      * Note that invoking this method does not itself necessarily transfer 
> any FlowFiles between connections. Rather, it informs the framework
>      * that the move needs to happen. It is up to the framework to make the 
> determination when this should happen. The framework may choose to move
>      * the FlowFiles inline during this method invocation, or it may choose 
> to do so later. The framework does guarantee, however, that this move
>      * happens before the Connector is able to start.
>      *
>      * <p>Any FlowFile that does not match the given predicate is ignored. 
> Many invocations may be used along with different Predicates in order to
>      * route FlowFiles from the given Connection to one of multiple 
> destination connections.
>      *
>      * @param removedConnectionId the ID of the Connection in the previous 
> flow
>      * @param destinationConnectionId the ID of the Connection in the current 
> flow
>      * @param filter a predicate that determines whether or not a given 
> FlowFile will be rehomed.
>      * @param attributesTransformer a transformer that can be used to 
> manipulate the attributes of FlowFiles that are moved
>      * @throws IllegalStateException if no Connection exists in the current 
> flow with the given ID,
>      *   or if a Connection still exists in the current flow with an ID 
> matching the removedConnectionId.
>      */
>     void relocateFlowFiles(String removedConnectionId, String 
> destinationConnectionId, Predicate<FlowFile> filter, AttributesTransformer 
> attributesTransformer);
>     /**
>      * Instructs the framework to permanently drop any FlowFiles that are 
> currently homed to the given removed Connection ID
>      *
>      * @param removedConnectionId the ID of the Connection in the previous 
> flow
>      * @throws IllegalStateException if a Connection still exists in the 
> current flow with an ID matching the removedConnectionId.
>      */
>     default void dropFlowFiles(String removedConnectionId) {
>         dropFlowFiles(removedConnectionId, flowFile -> true);
>     }
>     /**
>      * Instructs the framework to permanently drop any FlowFiles that are 
> currently homed to the given removed Connection ID if and only if
>      * they match the given predicate
>      *
>      * @param removedConnectionId the ID of the Connection in the previous 
> flow
>      * @param filter a predicate that determines whether or not a given 
> FlowFile will be dropped.
>      * @throws IllegalStateException if a Connection still exists in the 
> current flow with an ID matching the removedConnectionId.
>      */
>     void dropFlowFiles(String removedConnectionId, Predicate<FlowFile> 
> filter);
> }
> {code}
> And {{AttributesTransformer}} is a simple functional interface:
> {code:java}
> @FunctionalInterface
> public interface AttributesTransformer {
>     Map<String, String> transformAttributes(FlowFile flowFile);
>     AttributesTransformer NOP = FlowFile::getAttributes;
> } {code}
> The method on {{Connector}} would be defaulted to do nothing, so it is 
> backward compatible. But it would allow for scenarios such as "I used to have 
> Processor A -> B -> C -> D, but now I just want to do A -> X -> D" by 
> allowing the data from connection B-C and the data from C-D to be moved to 
> connection A-X:
> {code:java}
> public void evolveFlow(FlowEvolutionContext context) {
>   context.relocateFlowFiles( "B-C", "A-X" );
>   context.relocateFlowFiles( "C-D", "A-X", flowFile -> {
>     final Map<String, String> evolvedAttributes = new 
> HashMap<>(flowFile.getAttributes();
>     evolvedAttributes.remove("mime.type");
>     evolvedAttributes.put("attributeA", "valueA");
>     return evolvedAttributes;
>   });
> } {code}
> The {{Predicate}} allows us to also relocate different FlowFiles to different 
> Connections based on their attributes, size, etc.
> This approach means that we don't have to maintain old versions of the flow 
> or deal with starting up with that old version, draining, and then updating 
> the flow. Instead, we can simply have it move the data to a new Connection.
> As always, the intent of the NIP is to proposal the concept; the API may end 
> up deviating from this slightly, and will go through the normal 
> review-then-commit policy.
> h2. Provenance
> One important gotcha here is that the Provenance lineage could be very 
> confusing, even though it is technically correct in terms of what happened. 
> To account for this, we should introduce a new Provenance Event Type, such as 
> {{RELOCATED}} . This would make it much clearer what happened.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to