Pranathi A
03/26/2023, 10:58 AMMayank
Mayank
Pranathi A
03/26/2023, 2:17 PMMayank
francoisa
03/26/2023, 3:49 PMPadmini
12/10/2025, 12:25 PMPadmini
12/11/2025, 5:41 AMMayank
Shounak Kulkarni
12/11/2025, 6:51 AMRecordPurgerFactory or RecordModifierFactory being set in the minion starter, in that case I assume @Pranathi A /@Padmini would need to either extend the MinionStarter or update it to inject the purger into MinionContext right? https://github.com/apache/pinot/blob/master/pinot-plugins/pinot-minion-tasks/pinot[…]ache/pinot/plugin/minion/tasks/purge/PurgeTaskExecutorTest.javafrancoisa
12/11/2025, 1:15 PMpackage com.boondmanager.pinot;
import java.io.FileInputStream;
import java.io.IOException;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
import javax.net.ssl.SSLContext;
import org.apache.commons.io.FileUtils;
import org.apache.commons.lang.StringUtils;
import org.apache.helix.model.InstanceConfig;
import org.apache.helix.task.TaskStateModelFactory;
import org.apache.pinot.common.Utils;
import org.apache.pinot.common.auth.AuthProviderUtils;
import org.apache.pinot.common.config.TlsConfig;
import org.apache.pinot.common.metrics.MinionMeter;
import org.apache.pinot.common.metrics.MinionMetrics;
import org.apache.pinot.common.utils.ClientSSLContextGenerator;
import org.apache.pinot.common.utils.ServiceStatus;
import org.apache.pinot.common.utils.TlsUtils;
import org.apache.pinot.common.utils.fetcher.SegmentFetcherFactory;
import org.apache.pinot.common.utils.helix.HelixHelper;
import org.apache.pinot.core.util.ListenerConfigUtil;
import org.apache.pinot.minion.BaseMinionStarter;
import org.apache.pinot.minion.MinionAdminApiApplication;
import org.apache.pinot.minion.MinionContext;
import org.apache.pinot.minion.taskfactory.TaskFactoryRegistry;
import org.apache.pinot.spi.crypt.PinotCrypterFactory;
import org.apache.pinot.spi.env.PinotConfiguration;
import java.io.File;
import org.apache.pinot.spi.filesystem.PinotFSFactory;
import org.apache.pinot.spi.metrics.PinotMetricUtils;
import org.apache.pinot.spi.metrics.PinotMetricsRegistry;
import org.apache.pinot.spi.utils.CommonConstants;
import org.apache.pinot.sql.parsers.rewriter.QueryRewriterFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class CustomMinionStarter extends BaseMinionStarter {
private static final Logger LOGGER = LoggerFactory.getLogger(CustomMinionStarter.class);
private static final String HTTPS_ENABLED = "enabled";
public static void main(String[] args)
throws Exception {
if (args.length < 1) {
throw new IllegalArgumentException("Config file path is required");
}
String configFilePath = args[0];
Map<String, Object> configMap = new HashMap<>();
try (FileInputStream input = new FileInputStream(configFilePath)) {
Properties properties = new Properties();
properties.load(input);
// Convert Properties to Map<String, Object>
for (String name : properties.stringPropertyNames()) {
configMap.put(name, properties.getProperty(name));
}
} catch (IOException e) {
throw new RuntimeException("Failed to load configuration from file: " + configFilePath, e);
}
PinotConfiguration config = new PinotConfiguration(configMap);
CustomMinionStarter starter = new CustomMinionStarter();
starter.init(config);
starter.start();
}
@Override
public void start()
throws Exception {
LOGGER.error("Starting Pinot custom minion: {}", _instanceId);
Utils.logVersions();
MinionContext minionContext = MinionContext.getInstance();
minionContext.setRecordPurgerFactory(new OptimizedRecordPurgerFactory());
LOGGER.error("Record Purger registered in custom !");
// Initialize data directory
LOGGER.error("Initializing data directory");
File dataDir = new File(_config
.getProperty(CommonConstants.Helix.Instance.DATA_DIR_KEY, CommonConstants.Minion.DEFAULT_INSTANCE_DATA_DIR));
if (dataDir.exists()) {
FileUtils.cleanDirectory(dataDir);
} else {
FileUtils.forceMkdir(dataDir);
}
minionContext.setDataDir(dataDir);
// Initialize metrics
LOGGER.error("Initializing metrics");
// TODO: put all the metrics related configs down to "pinot.server.metrics"
PinotMetricsRegistry metricsRegistry = PinotMetricUtils.getPinotMetricsRegistry(_config.getMetricsConfig());
MinionMetrics minionMetrics = new MinionMetrics(_config.getMetricsPrefix(), metricsRegistry);
minionMetrics.initializeGlobalMeters();
minionContext.setMinionMetrics(minionMetrics);
// Install default SSL context if necessary (even if not force-enabled everywhere)
TlsConfig tlsDefaults = TlsUtils.extractTlsConfig(_config, CommonConstants.Minion.MINION_TLS_PREFIX);
if (StringUtils.isNotBlank(tlsDefaults.getKeyStorePath()) || StringUtils
.isNotBlank(tlsDefaults.getTrustStorePath())) {
<http://LOGGER.info|LOGGER.info>("Installing default SSL context for any client requests");
TlsUtils.installDefaultSSLSocketFactory(tlsDefaults);
}
// initialize authentication
minionContext.setTaskAuthProvider(
AuthProviderUtils.extractAuthProvider(_config, CommonConstants.Minion.CONFIG_TASK_AUTH_NAMESPACE));
// Start all components
<http://LOGGER.info|LOGGER.info>("Initializing PinotFSFactory");
PinotConfiguration pinotFSConfig = _config.subset(CommonConstants.Minion.PREFIX_OF_CONFIG_OF_PINOT_FS_FACTORY);
if (pinotFSConfig.isEmpty()) {
pinotFSConfig = _config.subset(CommonConstants.Minion.DEPRECATED_PREFIX_OF_CONFIG_OF_PINOT_FS_FACTORY);
}
PinotFSFactory.init(pinotFSConfig);
<http://LOGGER.info|LOGGER.info>("Initializing QueryRewriterFactory");
QueryRewriterFactory.init(_config.getProperty(CommonConstants.Minion.CONFIG_OF_MINION_QUERY_REWRITER_CLASS_NAMES));
<http://LOGGER.info|LOGGER.info>("Initializing segment fetchers for all protocols");
PinotConfiguration segmentFetcherFactoryConfig =
_config.subset(CommonConstants.Minion.PREFIX_OF_CONFIG_OF_SEGMENT_FETCHER_FACTORY);
if (segmentFetcherFactoryConfig.isEmpty()) {
segmentFetcherFactoryConfig =
_config.subset(CommonConstants.Minion.DEPRECATED_PREFIX_OF_CONFIG_OF_SEGMENT_FETCHER_FACTORY);
}
SegmentFetcherFactory.init(segmentFetcherFactoryConfig);
<http://LOGGER.info|LOGGER.info>("Initializing pinot crypter");
PinotConfiguration pinotCrypterConfig = _config.subset(CommonConstants.Minion.PREFIX_OF_CONFIG_OF_PINOT_CRYPTER);
if (pinotCrypterConfig.isEmpty()) {
pinotCrypterConfig = _config.subset(CommonConstants.Minion.DEPRECATED_PREFIX_OF_CONFIG_OF_PINOT_CRYPTER);
}
PinotCrypterFactory.init(pinotCrypterConfig);
// Need to do this before we start receiving state transitions.
<http://LOGGER.info|LOGGER.info>("Initializing ssl context for segment uploader");
PinotConfiguration segmentUploaderConfig =
_config.subset(CommonConstants.Minion.PREFIX_OF_CONFIG_OF_SEGMENT_UPLOADER);
if (segmentUploaderConfig.isEmpty()) {
segmentUploaderConfig =
_config.subset(CommonConstants.Minion.DEPRECATED_PREFIX_OF_CONFIG_OF_SEGMENT_UPLOADER);
}
PinotConfiguration httpsConfig = segmentUploaderConfig.subset(CommonConstants.HTTPS_PROTOCOL);
if (httpsConfig.getProperty(HTTPS_ENABLED, false)) {
SSLContext sslContext =
new ClientSSLContextGenerator(httpsConfig.subset(CommonConstants.PREFIX_OF_SSL_SUBSET)).generate();
minionContext.setSSLContext(sslContext);
}
// Join the Helix cluster
<http://LOGGER.info|LOGGER.info>("Joining the Helix cluster");
_helixManager.getStateMachineEngine().registerStateModelFactory("Task", new TaskStateModelFactory(_helixManager,
new TaskFactoryRegistry(_taskExecutorFactoryRegistry, _eventObserverFactoryRegistry).getTaskFactoryRegistry()));
_helixManager.connect();
updateInstanceConfigIfNeeded();
minionContext.setHelixPropertyStore(_helixManager.getHelixPropertyStore());
<http://LOGGER.info|LOGGER.info>("Starting minion admin application on: {}", ListenerConfigUtil.toString(_listenerConfigs));
_minionAdminApplication = new MinionAdminApiApplication(_instanceId, _config);
_minionAdminApplication.start(_listenerConfigs);
// Initialize health check callback
<http://LOGGER.info|LOGGER.info>("Initializing health check callback");
ServiceStatus.setServiceStatusCallback(_instanceId, new ServiceStatus.ServiceStatusCallback() {
private volatile boolean _isStarted = false;
private volatile String _statusDescription = "Helix ZK Not connected as " + _helixManager.getInstanceType();
@Override
public ServiceStatus.Status getServiceStatus() {
// TODO: add health check here
minionMetrics.addMeteredGlobalValue(MinionMeter.HEALTH_CHECK_GOOD_CALLS, 1L);
if (_isStarted) {
if (_helixManager.isConnected()) {
return ServiceStatus.Status.GOOD;
} else {
return ServiceStatus.Status.BAD;
}
}
if (!_helixManager.isConnected()) {
return ServiceStatus.Status.STARTING;
} else {
_isStarted = true;
_statusDescription = ServiceStatus.STATUS_DESCRIPTION_NONE;
return ServiceStatus.Status.GOOD;
}
}
@Override
public String getStatusDescription() {
return _statusDescription;
}
});
<http://LOGGER.info|LOGGER.info>("Pinot minion custom started");
}
private void updateInstanceConfigIfNeeded() {
InstanceConfig instanceConfig = HelixHelper.getInstanceConfig(_helixManager, _instanceId);
boolean updated = HelixHelper.updateHostnamePort(instanceConfig, _hostname, _port);
updated |= HelixHelper.addDefaultTags(instanceConfig,
() -> Collections.singletonList(CommonConstants.Helix.UNTAGGED_MINION_INSTANCE));
updated |= HelixHelper.removeDisabledPartitions(instanceConfig);
if (updated) {
HelixHelper.updateInstanceConfig(_helixManager, instanceConfig);
}
}
}francoisa
12/11/2025, 1:15 PMpackage com.boondmanager.pinot;
import java.io.BufferedReader;
import java.io.FileReader;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
import org.apache.pinot.core.minion.SegmentPurger;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class OptimizedRecordPurgerFactory implements SegmentPurger.RecordPurgerFactory {
private static final Logger LOGGER = LoggerFactory.getLogger(OptimizedRecordPurgerFactory.class);
// Charger et structurer les règles une seule fois
private Map<String, Map<String, Map<String, Set<String>>>> loadPurgeRules() {
Map<String, Map<String, Map<String, Set<String>>>> rules = new HashMap<>();
List<List<String>> listToPurge = getMapFromCSV(System.getenv("PURGE_LIST_FILE"));
if (listToPurge.isEmpty()) {
LOGGER.warn("No purge rules found");
return rules;
}
for (List<String> entry : listToPurge) {
if (entry.size() < 4) {
LOGGER.warn("Invalid entry format: {}", entry);
continue;
}
String customer = entry.get(0);
String table = entry.get(1);
String id = entry.get(2);
String ts = entry.get(3);
if (customer == null || table == null || id == null || ts == null) {
LOGGER.warn("Null values in entry: {}", entry);
continue;
}
rules.computeIfAbsent(table, k -> new HashMap<>())
.computeIfAbsent(customer, k -> new HashMap<>())
.computeIfAbsent(id, k-> new HashSet<>())
.add(ts);
}
return rules;
}
@Override
public SegmentPurger.RecordPurger getRecordPurger(String rawTableName) {
Map<String, Map<String, Map<String, Set<String>>>> purgeRules = loadPurgeRules();
// Optimisation : déterminer la table réelle une seule fois
String contextTable = getContextTable(rawTableName);
// Récupérer les règles de purge associées à cette table
Map<String, Map<String, Set<String>>> customerRules = purgeRules.containsKey("*") ? purgeRules.get("*") :
purgeRules.getOrDefault(contextTable, Collections.emptyMap());
if (customerRules == null || customerRules.isEmpty()) {
LOGGER.error("No suitable customerRules");
}
// Retourner un purger performant et sécurisé
return row -> {
// Vérification explicite du nom de la table
if (customerRules == null) {
return false;
}
String customer = String.valueOf(row.getValue("meta.customer"));
String id = String.valueOf(row.getValue("id"));
long ts = ((Number) row.getValue("meta.ts")).longValue();
Map<String, Set<String>> ids = customerRules.get(customer);
if (ids == null) return false;
String matchingID = getMatchingCondition(ids, id);
LOGGER.error(matchingID);
if (matchingID == null) {
return false;
}
return compareTimestamp(ids.get(matchingID),ts);
};
}
public String getMatchingCondition(Map<String, Set<String>> ids, String id) {
if (ids.containsKey("*")) {
return "*"; // Matched on "*"
}
if (ids.containsKey(id)) {
return id; // Matched on specific id
}
return null; // No match
}
private String getContextTable(String rawTableName) {
if (rawTableName.endsWith("_OFFLINE")) {
return rawTableName.substring(0, rawTableName.length() - 8);
} else if (rawTableName.endsWith("_REALTIME")) {
return rawTableName.substring(0, rawTableName.length() - 9);
}
return rawTableName;
}
public boolean compareTimestamp(Set<String> set, long value) {
return set.stream()
.findFirst()
.map(Long::parseLong)
.map(storedValue -> value <= storedValue)
.orElse(false);
}
public static List<List<String>> getMapFromCSV(final String filePath) {
List<List<String>> records = new ArrayList<>();
try (BufferedReader br = new BufferedReader(new FileReader(filePath))) {
String line;
while ((line = br.readLine()) != null) {
String[] values = line.split(Pattern.quote("|"));
records.add(Arrays.asList(values));
}
return records;
} catch (Exception e) {
LOGGER.error("Failed to parse purgeFile: {}", e.getMessage(), e);
return Collections.emptyList();
}
}francoisa
12/11/2025, 1:16 PMShounak Kulkarni
12/11/2025, 1:29 PMfrancoisa
12/11/2025, 1:31 PMPadmini
12/11/2025, 2:18 PMfrancoisa
12/11/2025, 2:29 PMPadmini
12/11/2025, 2:36 PMPadmini
12/11/2025, 2:38 PMfrancoisa
12/11/2025, 2:55 PMRetentionManager is not running that often. Have you tried to launch this periodictask by hand to see if it works ? Using the api on the swagger
or using curl
curl -X 'GET' \
'<http://controller-URL:9000/periodictask/run?taskname=RetentionManager&tableName=tableName&type=REALTIME>' \
-H 'accept: application/json' \
-H 'Authorization: Basic AUTHTOKEN'Padmini
12/11/2025, 3:14 PMPadmini
12/11/2025, 3:14 PMPadmini
12/11/2025, 3:21 PMfrancoisa
12/12/2025, 1:24 PMShounak Kulkarni
12/12/2025, 1:34 PMDeleting x segments from table: yfrancoisa
12/12/2025, 1:35 PMMayank
Padmini
12/16/2025, 7:42 AM