Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ public class NiFiProperties extends ApplicationProperties {
public static final String BORED_YIELD_DURATION = "nifi.bored.yield.duration";
public static final String SCHEDULING_STRATEGY = "nifi.scheduling.strategy";
public static final String PROCESSOR_SCHEDULING_TIMEOUT = "nifi.processor.scheduling.timeout";
public static final String PROCESSOR_AUTO_MAX_CONCURRENT_TASKS = "nifi.processor.auto.max.concurrent.tasks";
public static final String BACKPRESSURE_COUNT = "nifi.queue.backpressure.count";
public static final String BACKPRESSURE_SIZE = "nifi.queue.backpressure.size";
public static final String UPLOAD_WORKING_DIRECTORY = "nifi.upload.working.directory";
Expand Down Expand Up @@ -1477,6 +1478,22 @@ public String getSchedulingStrategy() {
return getProperty(SCHEDULING_STRATEGY, DEFAULT_SCHEDULING_STRATEGY);
}

public int getProcessorAutoMaxConcurrentTasks() {
final String configuredMaximum = getProperty(PROCESSOR_AUTO_MAX_CONCURRENT_TASKS, "12");
final int maximum;
try {
maximum = Integer.parseInt(configuredMaximum.trim());
} catch (final NumberFormatException e) {
throw new IllegalArgumentException(PROCESSOR_AUTO_MAX_CONCURRENT_TASKS + " must be a positive integer: " + configuredMaximum, e);
}

if (maximum < 1) {
throw new IllegalArgumentException(PROCESSOR_AUTO_MAX_CONCURRENT_TASKS + " must be a positive integer: " + configuredMaximum);
}

return (int) Math.min(maximum, 4L * Runtime.getRuntime().availableProcessors());
}

public File getStateManagementConfigFile() {
return new File(getProperty(STATE_MANAGEMENT_CONFIG_FILE, DEFAULT_STATE_MANAGEMENT_CONFIG_FILE));
}
Expand Down
1 change: 1 addition & 0 deletions nifi-docs/src/main/asciidoc/administration-guide.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -2953,6 +2953,7 @@ This cleanup mechanism takes into account only automatically created archived _f
|`nifi.administrative.yield.duration`|If a component allows an unexpected exception to escape, it is considered a bug. As a result, the framework will pause (or administratively yield) the component for this amount of time. This is done so that the component does not use up massive amounts of system resources, since it is known to have problems in the existing state. The default value is `30 secs`.
|`nifi.bored.yield.duration`|When a component has no work to do (i.e., is "bored"), this is the amount of time it will wait before checking to see if it has new data to work on. This way, it does not use up CPU resources by checking for new work too often. When setting this property, be aware that it could add extra latency for components that do not constantly have work to do, as once they go into this "bored" state, they will wait this amount of time before checking for more work. The default value is `10 ms`.
|`nifi.scheduling.strategy`|Selects the scheduling engine for Timer-Driven and Cron-Driven components. `AUTO` (the default) uses virtual threads on Java 25 or newer and standard scheduling on older Java versions. `VIRTUAL` always uses virtual threads. On Java 21, blocking inside synchronized component code can also block the underlying platform thread and reduce throughput. `STANDARD` uses a fixed platform thread pool sized according to the Maximum Timer Driven Thread Count configured in Controller Settings.
|`nifi.processor.auto.max.concurrent.tasks`|Sets the maximum number of concurrent tasks that a Processor using automatic scheduling is allowed to run. The default is `12` and the value must be a positive integer. NiFi also caps this value at four times the number of CPU cores available to the JVM. Automatic scheduling may use fewer concurrent tasks based on observed work and the instance's overall scheduling budget.
|`nifi.queue.backpressure.count`|When drawing a new connection between two components, this is the default value for that connection's back pressure object threshold. The default is `10000` and the value must be an integer.
|`nifi.queue.backpressure.size`|When drawing a new connection between two components, this is the default value for that connection's back pressure data size threshold. The default is `1 GB` and the value must be a data size including the unit of measure.
|`nifi.authorizer.configuration.file`*|This is the location of the file that specifies how authorizers are defined. The default value is `./conf/authorizers.xml`.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.nifi.jms.processors;

import jakarta.jms.Session;
import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.behavior.DynamicProperty;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
Expand Down Expand Up @@ -68,6 +69,8 @@
* properties that came with message which are added to a {@link FlowFile} as
* attributes.
*/
// Durable, non-shared subscribers require exactly one configured concurrent task.
@AllowsAutoScheduling(false)
@Tags({ "jms", "get", "message", "receive", "consume" })
@InputRequirement(Requirement.INPUT_FORBIDDEN)
@CapabilityDescription("Consumes JMS Message of type BytesMessage, TextMessage, ObjectMessage, MapMessage or StreamMessage transforming its content to "
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
*/
package org.apache.nifi.kafka.processors;

import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.WritesAttribute;
import org.apache.nifi.annotation.behavior.WritesAttributes;
Expand Down Expand Up @@ -131,6 +132,8 @@
})
@InputRequirement(InputRequirement.Requirement.INPUT_FORBIDDEN)
@SeeAlso({PublishKafka.class})
// Changing the number of concurrent tasks changes the number of Kafka consumers.
@AllowsAutoScheduling(false)
public class ConsumeKafka extends AbstractProcessor implements VerifiableProcessor, BacklogReportingProcessor {

static final AllowableValue TOPIC_NAME = new AllowableValue("names", "names", "Topic is a full topic name or comma separated list of names");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.nifi.processors.script;

import org.apache.commons.io.IOUtils;
import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.behavior.DynamicProperty;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
Expand Down Expand Up @@ -61,6 +62,8 @@
import javax.script.ScriptEngine;
import javax.script.SimpleBindings;

// Script engines are created eagerly for every ProcessContext concurrency slot.
@AllowsAutoScheduling(false)
@Tags({"script", "execute", "groovy", "clojure"})
@CapabilityDescription("Experimental - Executes a script given the flow file and a process session. The script is responsible for "
+ "handling the incoming flow file (transfer to SUCCESS or remove, e.g.) as well as any FlowFiles created by "
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.nifi.processors.script;

import org.apache.commons.io.IOUtils;
import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.lifecycle.OnScheduled;
import org.apache.nifi.components.PropertyDescriptor;
import org.apache.nifi.components.ValidationContext;
Expand Down Expand Up @@ -49,6 +50,8 @@
import javax.script.ScriptException;
import javax.script.SimpleBindings;

// Script engines are created eagerly for every ProcessContext concurrency slot.
@AllowsAutoScheduling(false)
abstract class ScriptedRecordProcessor extends AbstractProcessor implements Searchable {
protected static final Set<String> SCRIPT_OPTIONS = ScriptingComponentUtils.getAvailableEngines();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import io.krakens.grok.api.GrokCompiler;
import io.krakens.grok.api.Match;
import io.krakens.grok.api.exception.GrokException;
import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.SideEffectFree;
import org.apache.nifi.annotation.behavior.SupportsBatching;
Expand Down Expand Up @@ -66,6 +67,8 @@
@SupportsBatching
@SideEffectFree
@InputRequirement(InputRequirement.Requirement.INPUT_REQUIRED)
// A maximum-sized content buffer is allocated eagerly for every ProcessContext concurrency slot.
@AllowsAutoScheduling(false)
@Tags({"grok", "log", "text", "parse", "delimit", "extract"})
@CapabilityDescription("Evaluates one or more Grok Expressions against the content of a FlowFile, " +
"adding the results as attributes or replacing the content of the FlowFile with a JSON " +
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
*/
package org.apache.nifi.processors.standard;

import org.apache.nifi.annotation.behavior.AllowsAutoScheduling;
import org.apache.nifi.annotation.behavior.DynamicProperty;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
Expand Down Expand Up @@ -60,6 +61,8 @@
@SideEffectFree
@SupportsBatching
@InputRequirement(Requirement.INPUT_REQUIRED)
// A maximum-sized content buffer is allocated eagerly for every ProcessContext concurrency slot.
@AllowsAutoScheduling(false)
@Tags({"evaluate", "extract", "Text", "Regular Expression", "regex"})
@CapabilityDescription(
"Evaluates one or more Regular Expressions against the content of a FlowFile. "
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import org.apache.nifi.processor.FlowFileFilter;

import java.io.IOException;
import java.time.Instant;
import java.util.Collection;
import java.util.List;
import java.util.Set;
Expand All @@ -33,6 +34,25 @@

public interface FlowFileQueue {

/**
* Registers a listener for queue changes that can affect component readiness.
*
* @param listener listener to register
* @return registration used to remove the listener
*/
default QueueSchedulingRegistration addSchedulingListener(final QueueSchedulingListener listener) {
return QueueSchedulingRegistration.NO_OP;
}

/**
* Returns the next time at which the head FlowFile can become available without another queue mutation.
*
* @return the next availability time, or {@link Instant#EPOCH} when no deadline is known
*/
default Instant getNextFlowFileAvailabilityTime() {
return Instant.EPOCH;
}

/**
* @return the unique identifier for this FlowFileQueue
*/
Expand Down Expand Up @@ -97,6 +117,15 @@ public interface FlowFileQueue {

QueueSize size();

/**
* Returns work held in the local partition and available to this node.
*
* @return local queue size
*/
default QueueSize getLocalQueueSize() {
return getQueueDiagnostics().getLocalQueuePartitionDiagnostics().getActiveQueueSize();
}

/**
* Returns an atomic, point-in-time view of this queue's total {@link QueueSize} and the
* FlowFiles currently held in this node's active in-memory queue. Implementations freeze every
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
/*
* 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.nifi.controller.queue;

@FunctionalInterface
public interface QueueSchedulingListener {

/**
* Reports a queue change that can affect component readiness.
*/
void onQueueStateChanged();
}
Original file line number Diff line number Diff line change
@@ -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.
*/
package org.apache.nifi.controller.queue;

@FunctionalInterface
public interface QueueSchedulingRegistration extends AutoCloseable {

QueueSchedulingRegistration NO_OP = () -> {
};

/**
* Removes the registered queue scheduling listener. Repeated calls have no effect.
*/
@Override
void close();
}
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ public class ProcessorDTO extends ComponentDTO {
private String description;
private Boolean supportsParallelProcessing;
private Boolean supportsBatching;
private Boolean supportsAutoScheduling;
private Boolean supportsSensitiveDynamicProperties;
private Boolean supportsBacklogReporting;
private Boolean persistsState;
Expand Down Expand Up @@ -272,6 +273,18 @@ public void setSupportsBatching(Boolean supportsBatching) {
this.supportsBatching = supportsBatching;
}

/**
* @return whether this processor supports automatic scheduling
*/
@Schema(description = "Whether the processor supports automatic scheduling.")
public Boolean getSupportsAutoScheduling() {
return supportsAutoScheduling;
}

public void setSupportsAutoScheduling(final Boolean supportsAutoScheduling) {
this.supportsAutoScheduling = supportsAutoScheduling;
}

/**
* Gets the available relationships that this processor currently supports.
*
Expand Down
Loading
Loading