Skip to content
Merged
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 @@ -330,6 +330,10 @@ private IoTConsensusMessages() {}
public static final String LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ARG_ENTRY_S_LAST_ARG_MS_2AD8AB3D = "Subscription queue full, dropped {} entry(s) in the last {} ms, latest ";
public static final String LOG_SEARCHINDEX_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_2EA619ED = "searchIndex={}, queueSize={}, queueRemaining={}";
public static final String LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ENTRY_SEARCHINDEX_ARG_DROPPEDCOUNT_ARG_61F126B8 = "Subscription queue full, dropped entry searchIndex={}, droppedCount={}";
public static final String LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ARG_ENTRY_S_IN_THE_LAST_ARG_MS_WAL_REPLAY_REQUIRED_GROUP_ARG_LATEST_SEARCHINDEX_ARG_REASONCODE_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_0FBF6226 =
"Subscription realtime admission rejected {} entry(s) in the last {} ms; WAL replay required, group={}, latest searchIndex={}, reasonCode={}, queueSize={}, queueRemaining={}";
public static final String LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ENTRY_WAL_REPLAY_REQUIRED_GROUP_ARG_SEARCHINDEX_ARG_REASONCODE_ARG_REJECTEDCOUNT_ARG_7F76D6A9 =
"Subscription realtime admission rejected entry; WAL replay required, group={}, searchIndex={}, reasonCode={}, rejectedCount={}";
public static final String LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = "Reserved {} bytes for batch {}-{}, current total usage {}";
public static final String LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}: Failed to send idle writer safe-time barrier to {}. status={}";
public static final String LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576 = "{}: write operation failed. searchIndex: {}. Code: {}, subscriptionQueues: {}, this: {}";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,10 @@ private IoTConsensusMessages() {}
public static final String LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ARG_ENTRY_S_LAST_ARG_MS_2AD8AB3D = "订阅队列已满,丢弃了 {} 个 entry,最近 {} ms,最新 ";
public static final String LOG_SEARCHINDEX_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_2EA619ED = "searchIndex={}, queueSize={}, queueRemaining={}";
public static final String LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ENTRY_SEARCHINDEX_ARG_DROPPEDCOUNT_ARG_61F126B8 = "订阅队列已满,丢弃 entry,searchIndex={},droppedCount={}";
public static final String LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ARG_ENTRY_S_IN_THE_LAST_ARG_MS_WAL_REPLAY_REQUIRED_GROUP_ARG_LATEST_SEARCHINDEX_ARG_REASONCODE_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_0FBF6226 =
"订阅实时准入拒绝了 {} 个 entry(最近 {} ms);需通过 WAL 回放补齐,group={},最新 searchIndex={},reasonCode={},queueSize={},queueRemaining={}";
public static final String LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ENTRY_WAL_REPLAY_REQUIRED_GROUP_ARG_SEARCHINDEX_ARG_REASONCODE_ARG_REJECTEDCOUNT_ARG_7F76D6A9 =
"订阅实时准入拒绝 entry;需通过 WAL 回放补齐,group={},searchIndex={},reasonCode={},rejectedCount={}";
public static final String LOG_RESERVED_ARG_BYTES_BATCH_ARG_ARG_CURRENT_TOTAL_USAGE_ARG_308AE9C2 = "预留 {} 字节给批次 {}-{},当前总使用量 {}";
public static final String LOG_ARG_FAILED_SEND_IDLE_WRITER_SAFE_TIME_BARRIER_ARG_STATUS_AE047EAD = "{}:无法向 {} 发送 idle writer safe-time barrier。状态={}";
public static final String LOG_ARG_WRITE_OPERATION_FAILED_SEARCHINDEX_ARG_CODE_ARG_SUBSCRIPTIONQUEUES_ARG_THIS_ARG_F4B17576 = "{}:写入操作失败。searchIndex: {}。Code: {},订阅队列:{},当前对象:{}";
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
/*
* 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.iotdb.consensus.iot.subscription;

/** Exposes the reason for the calling thread's most recent subscription queue offer rejection. */
public interface SubscriptionQueueAdmission {

SubscriptionQueueRejectionReason getLastRejectionReason();
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,20 +32,22 @@
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;

public class SubscriptionQueueRegistry {

private static final Logger LOGGER = LoggerFactory.getLogger(SubscriptionQueueRegistry.class);

private static final long QUEUE_FULL_LOG_INTERVAL_MS = TimeUnit.SECONDS.toMillis(10);
private static final long REJECTION_LOG_INTERVAL_MS = TimeUnit.SECONDS.toMillis(10);

private final String consensusGroupId;
private final Map<BlockingQueue<IndexedConsensusRequest>, SubscriptionQueueRegistration> queues =
new ConcurrentHashMap<>();
private final AtomicLong droppedEntries = new AtomicLong();
private final AtomicLong lastDropLogTimeMs = new AtomicLong();
// offer is synchronized, so these counts remain scoped to one reason without extra atomics.
private final long[] rejectedEntriesByReason =
new long[SubscriptionQueueRejectionReason.values().length];
private final long[] lastRejectionLogTimeMsByReason =
new long[SubscriptionQueueRejectionReason.values().length];

public SubscriptionQueueRegistry(final String consensusGroupId) {
this.consensusGroupId = consensusGroupId;
Expand Down Expand Up @@ -141,27 +143,34 @@ public synchronized boolean offer(final IndexedConsensusRequest indexedConsensus
queue.remainingCapacity());
}
if (!offered) {
final long droppedCount = droppedEntries.incrementAndGet();
final SubscriptionQueueRejectionReason rejectionReason =
queue instanceof SubscriptionQueueAdmission
? ((SubscriptionQueueAdmission) queue).getLastRejectionReason()
: SubscriptionQueueRejectionReason.QUEUE_CAPACITY;
final int reasonIndex = rejectionReason.ordinal();
final long rejectedCount = ++rejectedEntriesByReason[reasonIndex];
final long now = System.currentTimeMillis();
final long lastLogTime = lastDropLogTimeMs.get();
if (now - lastLogTime >= QUEUE_FULL_LOG_INTERVAL_MS
&& lastDropLogTimeMs.compareAndSet(lastLogTime, now)) {
if (now - lastRejectionLogTimeMsByReason[reasonIndex] >= REJECTION_LOG_INTERVAL_MS) {
lastRejectionLogTimeMsByReason[reasonIndex] = now;
rejectedEntriesByReason[reasonIndex] = 0L;
LOGGER.warn(
IoTConsensusMessages
.LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ARG_ENTRY_S_LAST_ARG_MS_2AD8AB3D
+ IoTConsensusMessages
.LOG_SEARCHINDEX_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_2EA619ED,
droppedEntries.getAndSet(0),
QUEUE_FULL_LOG_INTERVAL_MS,
.LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ARG_ENTRY_S_IN_THE_LAST_ARG_MS_WAL_REPLAY_REQUIRED_GROUP_ARG_LATEST_SEARCHINDEX_ARG_REASONCODE_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_0FBF6226,
rejectedCount,
REJECTION_LOG_INTERVAL_MS,
consensusGroupId,
indexedConsensusRequest.getSearchIndex(),
rejectionReason.getCode(),
queue.size(),
queue.remainingCapacity());
} else if (LOGGER.isDebugEnabled()) {
LOGGER.debug(
IoTConsensusMessages
.LOG_SUBSCRIPTION_QUEUE_FULL_DROPPED_ENTRY_SEARCHINDEX_ARG_DROPPEDCOUNT_ARG_61F126B8,
.LOG_SUBSCRIPTION_REALTIME_ADMISSION_REJECTED_ENTRY_WAL_REPLAY_REQUIRED_GROUP_ARG_SEARCHINDEX_ARG_REASONCODE_ARG_REJECTEDCOUNT_ARG_7F76D6A9,
consensusGroupId,
indexedConsensusRequest.getSearchIndex(),
droppedCount);
rejectionReason.getCode(),
rejectedCount);
}
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
/*
* 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.iotdb.consensus.iot.subscription;

/** Stable error codes used when a consensus subscription queue rejects realtime input. */
public enum SubscriptionQueueRejectionReason {
NONE("NONE"),
QUEUE_CAPACITY("SUBSCRIPTION_QUEUE_CAPACITY"),
SUBSCRIPTION_MEMORY_QUOTA("SUBSCRIPTION_MEMORY_QUOTA"),
SUBSCRIPTION_MEMORY_LIMIT("SUBSCRIPTION_MEMORY_LIMIT"),
SUBSCRIPTION_OVERSIZED_ENTRY("SUBSCRIPTION_OVERSIZED_ENTRY"),
WRITER_BACKLOG("SUBSCRIPTION_WRITER_BACKLOG"),
SEEK_IN_PROGRESS("SUBSCRIPTION_SEEK_IN_PROGRESS"),
CONSENSUS_REQUEST_MEMORY_LIMIT("CONSENSUS_REQUEST_MEMORY_LIMIT"),
INACTIVE_OR_CLOSED("SUBSCRIPTION_QUEUE_INACTIVE_OR_CLOSED");

private final String code;

SubscriptionQueueRejectionReason(final String code) {
this.code = code;
}

public String getCode() {
return code;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2645,4 +2645,6 @@ private DataNodePipeMessages() {}
"Interrupted while reading OPC UA server operation limits, use defaults: maxNodesPerWrite={}, maxNodesPerNodeManagement={}";
public static final String LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871 =
"Failed to read OPC UA server operation limits, use defaults: maxNodesPerWrite={}, maxNodesPerNodeManagement={}";
public static final String MESSAGE_ARG_SUBSCRIPTION_ENTRY_REQUIRES_ARG_BYTES_ABOVE_THE_CURRENT_PER_QUEUE_MAXIMUM_ARG_BYTES_DATANODE_BUDGET_ARG_BYTES_REDUCE_THE_WRITE_BATCH_FIELD_SIZE_OR_INCREASE_SUBSCRIPTION_MATERIALIZATION_MEMORY_WAL_PROGRESS_HAS_NOT_ADVANCED_AFCBC7FC =
"[%s] Subscription entry requires %d bytes, above the current per-queue maximum %d bytes (DataNode budget %d bytes). Reduce the write batch/field size or increase subscription materialization memory. WAL progress has not advanced.";
}
Original file line number Diff line number Diff line change
Expand Up @@ -2469,4 +2469,6 @@ private DataNodePipeMessages() {}
"读取 OPC UA 服务器操作限制时被中断,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}";
public static final String LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871 =
"读取 OPC UA 服务器操作限制失败,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}";
public static final String MESSAGE_ARG_SUBSCRIPTION_ENTRY_REQUIRES_ARG_BYTES_ABOVE_THE_CURRENT_PER_QUEUE_MAXIMUM_ARG_BYTES_DATANODE_BUDGET_ARG_BYTES_REDUCE_THE_WRITE_BATCH_FIELD_SIZE_OR_INCREASE_SUBSCRIPTION_MATERIALIZATION_MEMORY_WAL_PROGRESS_HAS_NOT_ADVANCED_AFCBC7FC =
"[%s] 订阅 entry 需要 %d 字节,超过当前队列可用上限 %d 字节(DataNode 预算 %d 字节)。请减小写入 batch/字段大小或增加订阅物化内存。WAL 进度未推进。";
}
Loading
Loading