Hi Team, can someone suggest best approach for rea...
# pinot-dev
p
Hi Team, can someone suggest best approach for realtime table cleanup in Pinot?
m
What do you mean by cleanup? You want to delete a table, or bunch of tables, or segments of a table?
If just deleting a table, you can call delete api fro swagger.
p
Delete data in the table
m
Swagger should have api for that
f
Purgetask is the right way 😉
p
Hello, I want row based TTL on a particular table and I have set below configurations and Minion is also running. But it seems the purgeTask is not triggered because of below issues in minion.log Caught exception while executing task: Task_PurgeTask_2a125daa-00a1-4203-9e85-a8ddffa2332c_1765368600051_0 java.lang.IllegalArgumentException: At least one of record purger and modifier should be non-null controller.conf : controller.task.scheduler.enabled=true controller.minion.tasks.retention.enabled=true controller.retention.frequencyInSeconds=420 minion.conf : pinot.minion.task.scheduler.enabled=true pinot.minion.task.scheduler.frequencyInSeconds=60 pinot.minion.task.list=RetentionTask table_config.json task{ "PurgeTask":{ "schedule": "0 */5 * * * ?", "recordPurgerClassName": "org.apache.pinot.plugin.minion.tasks.purge.TimeBasedPurgerFactory", "recordPurgerConfigPrefix": "recordPurger", "recordPurger.timeColumnName": "insert_time", "recordPurger.retentionTimeValue": "5", "recordPurger.retentionTimeUnit": "MINUTES" } } Could anyone please suggest any solution.
Hi @Mayank /@francoisa can you please help
m
@Shounak Kulkarni ^^
s
Hey @francoisa, I don't see any default
RecordPurgerFactory
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.java
f
Here is the latest way I do what I expect for purging 😉 First step 2 classes 1- CustomMinionStarter witch register as RecordPurgerFactory my 2nd class
Copy code
package 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);
    }
  }


}
2- The core logic of my purging process 😉
Copy code
package 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();
    }
  }
Then I just manage to lauch my minion on the CustomMinionStarter 😉
s
Thanks Francoisa! Honestly I don't like the way we are using MinionContext to maintain the purger, its not easily pluggable, I rather have the decision made through task configs on which purger to pick than configuring something as part of starter and sticking with that. But for the above purpose if time based purge strategy is the one needed you can go with what Francoisa explained.
f
I honestly dont like it as well but no other documented ways 😇 If you want to delete records older than a given age there is also retention Policies on segement configs 😇
➕ 1
p
Hi @francoisa I tried retention policy as well but it didn't work. Could you provide any document to follow.
f
Share what you tried on table config 😇 It will be a better base to help out !
p
These are the config changes. table_config.json : "ingestionConfig": { "streamIngestionConfig": { "streamConfigMaps": [ { "streamType": "kafka", "realtime.segment.flush.threshold.time": "5m", "realtime.segment.flush.threshold.rows": "50", }}} "segmentsConfig": { "timeColumnName": "report_time", "timeType": "MILLISECONDS", "replication": "1", "retentionTimeValue": "5", "retentionTimeUnit": "MINUTES",} Controller.conf : controller.minion.enabled=true controller.minion.tasks.PurgeTask.enabled=true controller.task.scheduler.enabled=true controller.minion.tasks.retention.enabled=true controller.retention.frequencyInSeconds=420 Minion.conf : pinot.minion.enableTask=true pinot.minion.task.scheduler.enabled=true pinot.minion.task.scheduler.frequencyInSeconds=60 pinot.minion.task.list=retentionTask
@francoisa I have shared the config changes
f
Ok the
RetentionManager
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
Copy code
curl -X 'GET' \
  '<http://controller-URL:9000/periodictask/run?taskname=RetentionManager&tableName=tableName&type=REALTIME>' \
  -H 'accept: application/json' \
  -H 'Authorization: Basic AUTHTOKEN'
p
@francoisa RetentionManager is running but I see errors in controller.log
Start managing retention for table: demo_REALTIME 2025/12/11 202811.668 INFO [RetentionManager] [pool-12-thread-1] Segment lineage metadata clean-up is successfully processed for table: demo_REALTIME 2025/12/11 202811.668 INFO [RetentionManager] [pool-12-thread-1] Removing aged deleted segments for all tables 2025/12/11 202811.669 INFO [SegmentDeletionManager] [pool-12-thread-1] Fallback to using default cluster retention config: 604800000 ms 2025/12/11 202811.669 INFO [SegmentDeletionManager] [pool-12-thread-1] Fallback to using default cluster retention config: 604800000 ms 2025/12/11 202811.669 INFO [SegmentDeletionManager] [pool-12-thread-1] Fallback to using default cluster retention config: 604800000
Also the records are not getting deleted
f
Please do not DM me 😇 Open source means searching and try to figure out what happening. https://github.com/apache/pinot/blob/58e5c4afa8ef3f46f20c0a125469d40f5741ceab/pino[…]/apache/pinot/controller/helix/core/SegmentDeletionManager.java Coold be a good start to understood what’s going on
s
Hey Padmini can you share the segment metadata of the segments you think should be deleted but are not getting deleted? Asking this as a segment is only picked for deletion if all its records are beyond the retention boundary. Also check if the time column is configured properly (unit wise). If segmnets are picked for deletion you would see logs like
Deleting x segments from table: y
f
👌 1
m
Retention should not be used to include / excluded records. The best practice is to always include time filter in the query.
p
Hello , just a heads up . Sealed Segments are getting deleted now when I kept ingesting records continuously.
👍 2