pvillard31 commented on code in PR #11730:
URL: https://github.com/apache/nifi/pull/11730#discussion_r4206972933
##########
nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/scheduling/VirtualThreadSchedulingAgent.java:
##########
@@ -585,10 +1042,453 @@ private static OffsetDateTime getNextCronSchedule(final
OffsetDateTime currentSc
return cronExpression.next(now.isAfter(currentSchedule) ? now :
currentSchedule);
}
- private static class SchedulingGeneration {
+ private static class ProcessorTaskWorker {
+ private final long identifier;
+ private final AtomicBoolean retired = new AtomicBoolean();
+
+ private ProcessorTaskWorker(final long identifier) {
+ this.identifier = identifier;
+ }
+ }
+
+ private record WaitingProcessorTask(Thread thread, ProcessorTaskWorker
worker) {
+ }
+
+ private class ProcessorAutoSchedulingState implements
ConcurrentTaskMoveCoordinator.Participant {
+ private final Connectable connectable;
+ private final ConnectableTask connectableTask;
+ private final LifecycleState lifecycleState;
+ private final SchedulingGeneration generation;
+ private final int maxConcurrentTasks;
+ private final boolean batchingSupported;
+ private final boolean sourceProcessor;
+ private final StandardAutoSchedulingController controller;
+ private final ProcessorSchedulingMeasurements measurements = new
ProcessorSchedulingMeasurements();
+ private final AtomicReference<SchedulingSettings> schedulingSettings;
+ private final AtomicInteger activeProcessorInvocations = new
AtomicInteger();
+ private final ConcurrentMap<Long, ProcessorTaskWorker> workers = new
ConcurrentHashMap<>();
+ private final AtomicLong workerSequence = new AtomicLong();
+ private final AtomicBoolean sourceWorkCheckInProgress = new
AtomicBoolean();
+ private final AtomicLong resetSequence = new AtomicLong();
+ private final long evaluationIntervalNanos;
+ private final long baseSourceDelayNanos;
+ private final long maximumSourceDelayNanos;
+
+ private volatile long nextEvaluationNanos;
+ private volatile long lastEvaluationNanos = System.nanoTime();
+ private volatile long lastReadinessMeasurementNanos =
System.nanoTime();
+ private volatile long previousInputQueueSize;
+ private volatile long sourceDelayNanos;
+ private volatile long nextSourceWorkCheckNanos;
+ private volatile long lastSourceActivityNanos;
+ private volatile ProcessorSchedulingDecision lastDecision;
+ private volatile ProcessorSchedulingDecision lastChangeDecision;
+ private volatile ProcessorSchedulingSnapshot
lastProcessorSchedulingSnapshot;
+ private volatile boolean previousPrimaryNode;
+ private volatile boolean measurementsSupported;
+ private volatile double demandScore;
+
+ private ProcessorAutoSchedulingState(final Connectable connectable,
final ConnectableTask connectableTask,
+ final LifecycleState
lifecycleState, final SchedulingGeneration generation,
+ final int maxConcurrentTasks,
final boolean batchingSupported,
+ final boolean
processorTriggeredSerially) {
+ this.connectable = connectable;
+ this.connectableTask = connectableTask;
+ this.lifecycleState = lifecycleState;
+ this.generation = generation;
+ this.maxConcurrentTasks = maxConcurrentTasks;
+ this.batchingSupported = batchingSupported;
+ schedulingSettings = new AtomicReference<>(new
SchedulingSettings(1, batchingSupported ? TimeUnit.MILLISECONDS.toNanos(25L) :
0L));
+ sourceProcessor = connectable.isTriggerWhenEmpty() ||
!connectable.hasIncomingConnection() ||
!Connectables.hasNonLoopConnection(connectable);
+ controller = new
StandardAutoSchedulingController(maxConcurrentTasks, batchingSupported,
processorTriggeredSerially,
+ autoMaxConcurrentTasksBasedOnAvailableProcessors,
availableProcessorCount);
+ final double stagger =
((Math.floorMod(connectable.getIdentifier().hashCode(), 201) - 100) / 1000D);
+ this.evaluationIntervalNanos = (long)
(TimeUnit.SECONDS.toNanos(1L) * (1D + stagger));
+ this.nextEvaluationNanos = lastEvaluationNanos +
evaluationIntervalNanos;
+ this.baseSourceDelayNanos = Math.max(noWorkYieldNanos,
TimeUnit.MILLISECONDS.toNanos(1L));
+ this.maximumSourceDelayNanos = Math.max(baseSourceDelayNanos,
TimeUnit.MILLISECONDS.toNanos(100L));
+ this.previousPrimaryNode = flowController.isPrimary();
+ }
+
+ synchronized List<ProcessorTaskWorker> resizeWorkers(final int
targetWorkerCount) {
+ final List<ProcessorTaskWorker> activeWorkers = new ArrayList<>();
+ for (final ProcessorTaskWorker worker : workers.values()) {
+ if (!worker.retired.get()) {
+ activeWorkers.add(worker);
+ }
+ }
+
+ if (activeWorkers.size() > targetWorkerCount) {
+ activeWorkers.sort(Comparator.comparingLong(worker ->
worker.identifier));
+ for (int index = targetWorkerCount; index <
activeWorkers.size(); index++) {
+ activeWorkers.get(index).retired.set(true);
+ }
+
+ generation.signalChange();
+ generation.signalConcurrentTaskSlotChange();
+ return List.of();
+ }
+
+ final List<ProcessorTaskWorker> newWorkers = new ArrayList<>();
+ for (int index = activeWorkers.size(); index < targetWorkerCount;
index++) {
+ final ProcessorTaskWorker worker = new
ProcessorTaskWorker(workerSequence.getAndIncrement());
+ workers.put(worker.identifier, worker);
+ newWorkers.add(worker);
+ }
+
+ return newWorkers;
+ }
+
+ void workerStopped(final ProcessorTaskWorker worker) {
+ workers.remove(worker.identifier, worker);
+ generation.signalChange();
+ generation.signalConcurrentTaskSlotChange();
+ }
+
+ boolean tryAcquireConcurrentTaskSlot(final SchedulingSettings
expectedSettings) {
+ while (true) {
+ if (schedulingSettings.get() != expectedSettings) {
+ return false;
+ }
+
+ final int currentActiveProcessorInvocations =
activeProcessorInvocations.get();
+ if (currentActiveProcessorInvocations >=
expectedSettings.concurrentTasks()) {
+ return false;
+ }
+
+ if
(activeProcessorInvocations.compareAndSet(currentActiveProcessorInvocations,
currentActiveProcessorInvocations + 1)) {
+ if (schedulingSettings.get() == expectedSettings) {
+ return true;
+ }
+
+ activeProcessorInvocations.decrementAndGet();
+ generation.signalConcurrentTaskSlotChange();
+ return false;
+ }
+ }
+ }
+
+ void releaseConcurrentTaskSlot() {
+ activeProcessorInvocations.decrementAndGet();
+ generation.signalConcurrentTaskSlotChange();
+ }
+
+ boolean isSourceWorkCheckAllowed() {
+ return !sourceProcessor ||
connectableTask.hasLocallyConsumableInput() || sourceDelayNanos == 0L
+ || (System.nanoTime() >= nextSourceWorkCheckNanos &&
!sourceWorkCheckInProgress.get());
+ }
+
+ boolean reserveSourceWorkCheck() {
+ return requiresSourceWorkCheck() &&
sourceWorkCheckInProgress.compareAndSet(false, true);
+ }
+
+ boolean requiresSourceWorkCheck() {
+ return sourceProcessor &&
!connectableTask.hasLocallyConsumableInput() && sourceDelayNanos > 0L;
+ }
+
+ void releaseSourceWorkCheck() {
+ sourceWorkCheckInProgress.set(false);
+ generation.signalChange();
+ }
+
+ long getSourceWorkCheckDelayNanos() {
+ return Math.max(0L, nextSourceWorkCheckNanos - System.nanoTime());
+ }
+
+ void recordInvocationResult(final InvocationResult result) {
+ if (!sourceProcessor) {
+ return;
+ }
+
+ if (result.getOutcome() ==
InvocationOutcome.INVOKED_WITH_ACTIVITY) {
+ sourceDelayNanos = 0L;
+ nextSourceWorkCheckNanos = 0L;
+ lastSourceActivityNanos = System.nanoTime();
+ generation.signalChange();
+ } else if (result.getOutcome() ==
InvocationOutcome.INVOKED_WITHOUT_ACTIVITY
+ && activeProcessorInvocations.get() <= 1
+ && !connectableTask.hasLocallyConsumableInput()) {
+ sourceDelayNanos = sourceDelayNanos == 0L ?
baseSourceDelayNanos : Math.min(maximumSourceDelayNanos, sourceDelayNanos * 2L);
+ nextSourceWorkCheckNanos = System.nanoTime() +
sourceDelayNanos;
+ }
+ }
+
+ Instant getNextQueueDeadline() {
+ Instant earliestDeadline = Instant.EPOCH;
+ for (final FlowFileQueue queue : generation.queues) {
+ final Instant deadline =
queue.getNextFlowFileAvailabilityTime();
+ if (deadline.isAfter(Instant.EPOCH) &&
(earliestDeadline.equals(Instant.EPOCH) ||
deadline.isBefore(earliestDeadline))) {
+ earliestDeadline = deadline;
+ }
+ }
+
+ return earliestDeadline;
+ }
+
+ boolean isEvaluationDue(final long nowNanos) {
+ if (nowNanos < nextEvaluationNanos) {
+ return false;
+ }
+
+ nextEvaluationNanos = nowNanos + evaluationIntervalNanos;
+ return true;
+ }
+
+ void recordReadiness(final long nowNanos) {
+ final long measurementNanos = Math.max(1L, nowNanos -
lastReadinessMeasurementNanos);
+ lastReadinessMeasurementNanos = nowNanos;
+
measurements.recordReadiness(connectableTask.getReadinessOutcome(),
measurementNanos);
+ }
+
+ ProcessorSchedulingSnapshot captureProcessorSchedulingSnapshot(final
long nowNanos, final boolean globalCapacityTestAllowed) {
+ final SchedulingSettings currentSettings =
schedulingSettings.get();
+ long inputQueueSize = 0L;
+ for (final Connection connection :
connectable.getIncomingConnections()) {
Review Comment:
getLocalQueueSize() scans the active queue and reads every swap file
summary, but excludes swapped FlowFiles from the returned count. Could we use a
cheap queue-size counter that includes the full local backlog?
##########
nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/scheduling/VirtualThreadSchedulingAgent.java:
##########
@@ -54,23 +92,56 @@ public class VirtualThreadSchedulingAgent implements
SchedulingAgent {
private static final Logger logger =
LoggerFactory.getLogger(VirtualThreadSchedulingAgent.class);
private static final long PERMIT_POLL_INTERVAL_NANOS =
TimeUnit.SECONDS.toNanos(1L);
+ private static final long NON_BATCHED_PROCESSOR_BURST_NANOS =
TimeUnit.MILLISECONDS.toNanos(10L);
private final FlowController flowController;
private final RepositoryContextFactory contextFactory;
private final DynamicSemaphore globalSemaphore;
+ private final int autoMaxConcurrentTasks;
+ private final int availableProcessorCount;
+ private final boolean autoMaxConcurrentTasksBasedOnAvailableProcessors;
private final long noWorkYieldNanos;
private final ExecutorService executorService;
+ private final ScheduledExecutorService schedulingEvaluationExecutor;
+ private final SystemSchedulingMetrics systemSchedulingMetrics;
private final ConcurrentMap<String, SchedulingGeneration>
schedulingGenerations = new ConcurrentHashMap<>();
+ private final ConcurrentTaskMoveCoordinator concurrentTaskMoveCoordinator
= new ConcurrentTaskMoveCoordinator();
+ private final LongAdder globalPermitHoldNanos = new LongAdder();
+ private final LongAdder globalPermitWaitNanos = new LongAdder();
+ private final LongAdder globalTaskSlotsFullyUsedSamples = new LongAdder();
+ private final LongAdder globalTaskUsageSamples = new LongAdder();
private final AtomicBoolean shutdown = new AtomicBoolean();
private final AtomicInteger runningThreadCount = new AtomicInteger();
+ private volatile SystemSchedulingSnapshot systemSchedulingSnapshot =
SystemSchedulingSnapshot.createBuilder()
+ .setCpuLoad(-1D)
+ .setAverageCpuLoad(-1D)
+ .setMaxGlobalConcurrentTasks(1)
+ .build();
+ private volatile long lastSystemSchedulingSnapshotNanos =
System.nanoTime();
private volatile String adminYieldDuration = "1 sec";
private volatile long adminYieldNanos = TimeUnit.SECONDS.toNanos(1L);
public VirtualThreadSchedulingAgent(final FlowController flowController,
final RepositoryContextFactory contextFactory,
final NiFiProperties nifiProperties,
final int maxThreadCount) {
+ this(flowController, contextFactory, nifiProperties, maxThreadCount,
nifiProperties.getProcessorAutoMaxConcurrentTasks());
+ }
+
+ public VirtualThreadSchedulingAgent(final FlowController flowController,
final RepositoryContextFactory contextFactory,
+ final NiFiProperties nifiProperties,
final int maxThreadCount, final int autoMaxConcurrentTasks) {
+ this(flowController, contextFactory, nifiProperties, maxThreadCount,
autoMaxConcurrentTasks, new SystemSchedulingMetrics());
+ }
+
+ VirtualThreadSchedulingAgent(final FlowController flowController, final
RepositoryContextFactory contextFactory, final NiFiProperties nifiProperties,
+ final int maxThreadCount, final int
autoMaxConcurrentTasks, final SystemSchedulingMetrics systemSchedulingMetrics) {
+ this.systemSchedulingMetrics = systemSchedulingMetrics;
this.flowController = flowController;
this.contextFactory = contextFactory;
this.globalSemaphore = new DynamicSemaphore(maxThreadCount);
+ this.autoMaxConcurrentTasks = autoMaxConcurrentTasks;
+ availableProcessorCount = Runtime.getRuntime().availableProcessors();
+ final String configuredMaximum =
nifiProperties.getProperty(NiFiProperties.PROCESSOR_AUTO_MAX_CONCURRENT_TASKS);
Review Comment:
If this property is blank, can Integer.parseInt() fail during startup even
though NiFiProperties already falls back to the default value?
##########
nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/queue/clustered/SocketLoadBalancedFlowFileQueue.java:
##########
@@ -1312,39 +1373,36 @@ public void onLocalNodeIdentifierSet(final
NodeIdentifier localNodeId) {
} finally {
partitionWriteLock.unlock();
}
+
+ notifySchedulingListeners();
}
@Override
public void onNodeStateChange(final NodeIdentifier nodeId, final
NodeConnectionState newState) {
- partitionWriteLock.lock();
- try {
- if (!offloaded) {
- switch (newState) {
- case OFFLOADING:
- onNodeRemoved(nodeId);
- break;
- case CONNECTED:
- onNodeAdded(nodeId);
- break;
- }
- } else {
- switch (newState) {
- case CONNECTED:
- if (nodeId != null &&
nodeId.equals(clusterCoordinator.getLocalNodeIdentifier())) {
- // the node with this queue was connected to
the cluster, make sure the queue is not offloaded
- resetOffloadedQueue();
- }
- break;
- case OFFLOADED:
- case OFFLOADING:
- case DISCONNECTED:
- case DISCONNECTING:
- onNodeRemoved(nodeId);
- break;
- }
+ if (!offloaded) {
Review Comment:
Was removing partitionWriteLock around this state transition intentional,
given that resetOffloadedQueue() can now run concurrently with offloadQueue()?
##########
nifi-framework-bundle/nifi-framework/nifi-framework-components/src/main/java/org/apache/nifi/groups/StandardProcessGroup.java:
##########
@@ -1504,6 +1504,11 @@ public void addConnection(final Connection connection) {
} finally {
writeLock.unlock();
}
+
+ scheduler.notifySchedulingEvent(connection.getSource());
Review Comment:
Should changing a connection destination also notify the scheduler so a
running automatic Processor registers a listener on the moved queue?
##########
nifi-framework-bundle/nifi-framework/nifi-framework-core/src/main/java/org/apache/nifi/controller/scheduling/VirtualThreadSchedulingAgent.java:
##########
@@ -158,6 +263,155 @@ public void schedule(final Connectable connectable, final
LifecycleState lifecyc
}
}
+ private void registerQueueListeners(final Connectable connectable, final
SchedulingGeneration generation) {
+ final Set<FlowFileQueue> queues = new HashSet<>();
+ for (final Connection connection :
connectable.getIncomingConnections()) {
+ queues.add(connection.getFlowFileQueue());
+ }
+
+ for (final Connection connection : connectable.getConnections()) {
+ queues.add(connection.getFlowFileQueue());
+ }
+
+ for (final FlowFileQueue queue : queues) {
+ generation.addQueue(queue);
+
generation.addQueueRegistration(queue.addSchedulingListener(generation::signalChange));
+ }
+ }
+
+ private void startAutoWorkers(final ProcessorAutoSchedulingState state,
final int desiredWorkerCount) {
+ final List<ProcessorTaskWorker> workers =
state.resizeWorkers(desiredWorkerCount);
+ for (final ProcessorTaskWorker worker : workers) {
+ final String threadName = buildThreadName(state.connectable,
Math.toIntExact(worker.identifier));
+ submitTask(threadName, state.generation, () ->
runAutoSchedulingLoop(state, worker));
+ }
+ }
+
+ private void runAutoSchedulingLoop(final ProcessorAutoSchedulingState
state, final ProcessorTaskWorker worker) {
Review Comment:
Should this loop catch and log unexpected failures and apply the
administrative yield, as the existing scheduling loop does?
##########
nifi-extension-bundles/nifi-standard-bundle/nifi-standard-processors/src/main/java/org/apache/nifi/processors/standard/ExtractText.java:
##########
@@ -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)
Review Comment:
Should ListenUDP and ListenUDPRecord also opt out of automatic scheduling
because they eagerly allocate one buffer per concurrency slot?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]