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
15 changes: 15 additions & 0 deletions docs/docs/maintenance/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,11 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca
<td>Gauge</td>
<td>Number of scanned manifest files in the last scan.</td>
</tr>
<tr>
<td>lastScanSkippedManifests</td>
<td>Gauge</td>
<td>Number of manifest files skipped by manifest level filtering in the last scan.</td>
</tr>
<tr>
<td>lastScanSkippedTableFiles</td>
<td>Gauge</td>
Expand All @@ -80,6 +85,16 @@ Below is lists of Paimon built-in metrics. They are summarized into types of sca
<td>Gauge</td>
<td>Resulted table files in the last scan.</td>
</tr>
<tr>
<td>lastScanResultedTableFilesSize</td>
<td>Gauge</td>
<td>Total size in bytes of the resulted table files to be read in the last scan.</td>
</tr>
<tr>
<td>lastScanResultedRecordCount</td>
<td>Gauge</td>
<td>Total number of records in the resulted table files to be read in the last scan.</td>
</tr>
</tbody>
</table>

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import org.apache.paimon.manifest.BucketFilter;
import org.apache.paimon.manifest.FileEntry;
import org.apache.paimon.manifest.FileEntry.Identifier;
import org.apache.paimon.manifest.FileKind;
import org.apache.paimon.manifest.ManifestEntry;
import org.apache.paimon.manifest.ManifestEntrySerializer;
import org.apache.paimon.manifest.ManifestFile;
Expand Down Expand Up @@ -319,13 +320,26 @@ public Plan plan() {
manifestsResult.allManifests.stream()
.mapToLong(f -> f.numAddedFiles() - f.numDeletedFiles())
.sum();
// for DELTA and CHANGELOG scan modes the result contains both ADD and DELETE entries,
// only ADD entries will actually be read, so size and record count only count them
long resultedTableFilesSize = 0L;
long resultedRecordCount = 0L;
for (ManifestEntry entry : result) {
if (entry.kind() == FileKind.ADD) {
resultedTableFilesSize += entry.file().fileSize();
resultedRecordCount += entry.file().rowCount();
}
}
scanMetrics.reportScan(
new ScanStats(
scanDuration,
snapshot == null ? 0 : snapshot.id(),
manifests.size(),
manifestsResult.allManifests.size() - manifests.size(),
allDataFiles - result.size(),
result.size()));
result.size(),
resultedTableFilesSize,
resultedRecordCount));
}

return new Plan() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,12 @@ public class ScanMetrics {
public static final String SCAN_DURATION = "scanDuration";
public static final String LAST_SCANNED_SNAPSHOT_ID = "lastScannedSnapshotId";
public static final String LAST_SCANNED_MANIFESTS = "lastScannedManifests";
public static final String LAST_SCAN_SKIPPED_MANIFESTS = "lastScanSkippedManifests";
public static final String LAST_SCAN_SKIPPED_TABLE_FILES = "lastScanSkippedTableFiles";
public static final String LAST_SCAN_RESULTED_TABLE_FILES = "lastScanResultedTableFiles";
public static final String LAST_SCAN_RESULTED_TABLE_FILES_SIZE =
"lastScanResultedTableFilesSize";
public static final String LAST_SCAN_RESULTED_RECORD_COUNT = "lastScanResultedRecordCount";
public static final String MANIFEST_HIT_CACHE = "manifestHitCache";
public static final String MANIFEST_MISSED_CACHE = "manifestMissedCache";
public static final String DVMETA_HIT_CACHE = "dvMetaHitCache";
Expand All @@ -59,12 +63,21 @@ public ScanMetrics(MetricRegistry registry, String tableName) {
metricGroup.gauge(
LAST_SCANNED_MANIFESTS,
() -> latestScan == null ? 0L : latestScan.getScannedManifests());
metricGroup.gauge(
LAST_SCAN_SKIPPED_MANIFESTS,
() -> latestScan == null ? 0L : latestScan.getSkippedManifests());
metricGroup.gauge(
LAST_SCAN_SKIPPED_TABLE_FILES,
() -> latestScan == null ? 0L : latestScan.getSkippedTableFiles());
metricGroup.gauge(
LAST_SCAN_RESULTED_TABLE_FILES,
() -> latestScan == null ? 0L : latestScan.getResultedTableFiles());
metricGroup.gauge(
LAST_SCAN_RESULTED_TABLE_FILES_SIZE,
() -> latestScan == null ? 0L : latestScan.getResultedTableFilesSize());
metricGroup.gauge(
LAST_SCAN_RESULTED_RECORD_COUNT,
() -> latestScan == null ? 0L : latestScan.getResultedRecordCount());
metricGroup.gauge(MANIFEST_HIT_CACHE, () -> cacheMetrics.getHitObject().get());
metricGroup.gauge(MANIFEST_MISSED_CACHE, () -> cacheMetrics.getMissedObject().get());
metricGroup.gauge(DVMETA_HIT_CACHE, () -> dvMetaCacheMetrics.getHitObject().get());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,20 +26,30 @@ public class ScanStats {
private final long duration;
private final long scannedSnapshotId;
private final long scannedManifests;
private final long skippedManifests;
private final long skippedTableFiles;
private final long resultedTableFiles;
// the unit is bytes
private final long resultedTableFilesSize;
private final long resultedRecordCount;

public ScanStats(
long duration,
long scannedSnapshotId,
long scannedManifests,
long skippedManifests,
long skippedTableFiles,
long resultedTableFiles) {
long resultedTableFiles,
long resultedTableFilesSize,
long resultedRecordCount) {
this.duration = duration;
this.scannedSnapshotId = scannedSnapshotId;
this.scannedManifests = scannedManifests;
this.skippedManifests = skippedManifests;
this.skippedTableFiles = skippedTableFiles;
this.resultedTableFiles = resultedTableFiles;
this.resultedTableFilesSize = resultedTableFilesSize;
this.resultedRecordCount = resultedRecordCount;
}

@VisibleForTesting
Expand All @@ -52,6 +62,11 @@ protected long getScannedManifests() {
return scannedManifests;
}

@VisibleForTesting
protected long getSkippedManifests() {
return skippedManifests;
}

@VisibleForTesting
protected long getSkippedTableFiles() {
return skippedTableFiles;
Expand All @@ -62,6 +77,16 @@ protected long getResultedTableFiles() {
return resultedTableFiles;
}

@VisibleForTesting
protected long getResultedTableFilesSize() {
return resultedTableFilesSize;
}

@VisibleForTesting
protected long getResultedRecordCount() {
return resultedRecordCount;
}

@VisibleForTesting
protected long getDuration() {
return duration;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,11 @@ public void testGenericMetricsRegistration() {
ScanMetrics.SCAN_DURATION,
ScanMetrics.LAST_SCANNED_SNAPSHOT_ID,
ScanMetrics.LAST_SCANNED_MANIFESTS,
ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS,
ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES,
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES,
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE,
ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT,
ScanMetrics.MANIFEST_HIT_CACHE,
ScanMetrics.MANIFEST_MISSED_CACHE,
ScanMetrics.DVMETA_HIT_CACHE,
Expand All @@ -71,20 +74,32 @@ public void testMetricsAreUpdated() {
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_SNAPSHOT_ID);
Gauge<Long> lastScannedManifests =
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCANNED_MANIFESTS);
Gauge<Long> lastScanSkippedManifests =
(Gauge<Long>) registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS);
Gauge<Long> lastScanSkippedTableFiles =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES);
Gauge<Long> lastScanResultedTableFiles =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES);
Gauge<Long> lastScanResultedTableFilesSize =
(Gauge<Long>)
registeredGenericMetrics.get(
ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE);
Gauge<Long> lastScanResultedRecordCount =
(Gauge<Long>)
registeredGenericMetrics.get(ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT);

assertThat(lastScanDuration.getValue()).isEqualTo(0);
assertThat(lastScannedSnapshotId.getValue()).isEqualTo(0);
assertThat(scanDuration.getCount()).isEqualTo(0);
assertThat(scanDuration.getStatistics().size()).isEqualTo(0);
assertThat(lastScannedManifests.getValue()).isEqualTo(0);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(0);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(0);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(0);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(0);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(0);

// report once
reportOnce(scanMetrics);
Expand All @@ -101,8 +116,11 @@ public void testMetricsAreUpdated() {
assertThat(scanDuration.getStatistics().getMax()).isEqualTo(200);
assertThat(scanDuration.getStatistics().getStdDev()).isEqualTo(0);
assertThat(lastScannedManifests.getValue()).isEqualTo(20);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(5);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(25);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(10);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(1024);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(100);

// report again
reportAgain(scanMetrics);
Expand All @@ -119,17 +137,20 @@ public void testMetricsAreUpdated() {
assertThat(scanDuration.getStatistics().getMax()).isEqualTo(500);
assertThat(scanDuration.getStatistics().getStdDev()).isCloseTo(212.132, offset(0.001));
assertThat(lastScannedManifests.getValue()).isEqualTo(22);
assertThat(lastScanSkippedManifests.getValue()).isEqualTo(7);
assertThat(lastScanSkippedTableFiles.getValue()).isEqualTo(30);
assertThat(lastScanResultedTableFiles.getValue()).isEqualTo(8);
assertThat(lastScanResultedTableFilesSize.getValue()).isEqualTo(2048);
assertThat(lastScanResultedRecordCount.getValue()).isEqualTo(200);
}

private void reportOnce(ScanMetrics scanMetrics) {
ScanStats scanStats = new ScanStats(200, 1L, 20, 25, 10);
ScanStats scanStats = new ScanStats(200, 1L, 20, 5, 25, 10, 1024, 100);
scanMetrics.reportScan(scanStats);
}

private void reportAgain(ScanMetrics scanMetrics) {
ScanStats scanStats = new ScanStats(500, 2L, 22, 30, 8);
ScanStats scanStats = new ScanStats(500, 2L, 22, 7, 30, 8, 2048, 200);
scanMetrics.reportScan(scanStats);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,10 @@ abstract class PaimonBaseScan(table: InnerTable)
PaimonPlanningDurationMetric(),
PaimonScannedSnapshotIdMetric(),
PaimonScannedManifestsMetric(),
PaimonSkippedTableFilesMetric()
PaimonSkippedManifestsMetric(),
PaimonSkippedTableFilesMetric(),
PaimonResultedTableFilesSizeMetric(),
PaimonResultedRecordCountMetric()
)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,8 +29,11 @@ object PaimonMetrics {
val PLANNING_DURATION = "planningDuration"
val SCANNED_SNAPSHOT_ID = "scannedSnapshotId"
val SCANNED_MANIFESTS = "scannedManifests"
val SKIPPED_MANIFESTS = "skippedManifests"
val SKIPPED_TABLE_FILES = "skippedTableFiles"
val RESULTED_TABLE_FILES = "resultedTableFiles"
val RESULTED_TABLE_FILES_SIZE = "resultedTableFilesSize"
val RESULTED_RECORD_COUNT = "resultedRecordCount"
val RESULTED_POSTPONE_FILES = "resultedPostponeFiles"
val NUM_POSTPONE_RECORDS = "numPostponeRecords"

Expand Down Expand Up @@ -125,7 +128,7 @@ case class PaimonReadBatchTimeTaskMetric(value: Long) extends PaimonTaskMetric {

case class PaimonPlanningDurationMetric() extends PaimonTimingSumMetric {
override def name(): String = PaimonMetrics.PLANNING_DURATION
override def description(): String = "planing duration"
override def description(): String = "planning duration"
}

case class PaimonPlanningDurationTaskMetric(value: Long) extends PaimonTaskMetric {
Expand All @@ -150,6 +153,15 @@ case class PaimonScannedManifestsTaskMetric(value: Long) extends PaimonTaskMetri
override def name(): String = PaimonMetrics.SCANNED_MANIFESTS
}

case class PaimonSkippedManifestsMetric() extends PaimonSumMetric {
override def name(): String = PaimonMetrics.SKIPPED_MANIFESTS
override def description(): String = "number of skipped manifests"
}

case class PaimonSkippedManifestsTaskMetric(value: Long) extends PaimonTaskMetric {
override def name(): String = PaimonMetrics.SKIPPED_MANIFESTS
}

case class PaimonSkippedTableFilesMetric() extends PaimonSumMetric {
override def name(): String = PaimonMetrics.SKIPPED_TABLE_FILES
override def description(): String = "number of skipped table files"
Expand All @@ -168,6 +180,24 @@ case class PaimonResultedTableFilesTaskMetric(value: Long) extends PaimonTaskMet
override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES
}

case class PaimonResultedTableFilesSizeMetric() extends PaimonSizeSumMetric {
override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES_SIZE
override def description(): String = "size of resulted table files"
}

case class PaimonResultedTableFilesSizeTaskMetric(value: Long) extends PaimonTaskMetric {
override def name(): String = PaimonMetrics.RESULTED_TABLE_FILES_SIZE
}

case class PaimonResultedRecordCountMetric() extends PaimonSumMetric {
override def name(): String = PaimonMetrics.RESULTED_RECORD_COUNT
override def description(): String = "number of resulted records"
}

case class PaimonResultedRecordCountTaskMetric(value: Long) extends PaimonTaskMetric {
override def name(): String = PaimonMetrics.RESULTED_RECORD_COUNT
}

case class PaimonResultedPostponeFilesMetric() extends PaimonSumMetric {
override def name(): String = PaimonMetrics.RESULTED_POSTPONE_FILES
override def description(): String = "number of resulted postpone files"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,10 +50,16 @@ case class SparkMetricRegistry() extends MetricRegistry {
gauge[Long](metrics, ScanMetrics.LAST_SCANNED_SNAPSHOT_ID)),
PaimonScannedManifestsTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCANNED_MANIFESTS)),
PaimonSkippedManifestsTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCAN_SKIPPED_MANIFESTS)),
PaimonSkippedTableFilesTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCAN_SKIPPED_TABLE_FILES)),
PaimonResultedTableFilesTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES))
gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES)),
PaimonResultedTableFilesSizeTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_TABLE_FILES_SIZE)),
PaimonResultedRecordCountTaskMetric(
gauge[Long](metrics, ScanMetrics.LAST_SCAN_RESULTED_RECORD_COUNT))
)
case None =>
Array.empty
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@

package org.apache.paimon.spark.sql

import org.apache.paimon.spark.PaimonMetrics.{RESULTED_TABLE_FILES, SCANNED_SNAPSHOT_ID, SKIPPED_TABLE_FILES}
import org.apache.paimon.spark.PaimonMetrics.{RESULTED_RECORD_COUNT, RESULTED_TABLE_FILES, RESULTED_TABLE_FILES_SIZE, SCANNED_SNAPSHOT_ID, SKIPPED_MANIFESTS, SKIPPED_TABLE_FILES}
import org.apache.paimon.spark.PaimonSparkTestBase
import org.apache.paimon.spark.read.PaimonSplitScan
import org.apache.paimon.spark.util.ScanPlanHelper
Expand Down Expand Up @@ -53,25 +53,33 @@ class PaimonMetricTest extends PaimonSparkTestBase with ScanPlanHelper {
s: String,
scannedSnapshotId: Long,
skippedTableFiles: Long,
resultedTableFiles: Long): Unit = {
resultedTableFiles: Long,
resultedRecordCount: Long): Unit = {
val scan = getPaimonScan(s)
// call getInputPartitions to trigger scan
scan.inputPartitions
val metrics = scan.reportDriverMetrics()
Assertions.assertEquals(scannedSnapshotId, metric(metrics, SCANNED_SNAPSHOT_ID))
Assertions.assertEquals(skippedTableFiles, metric(metrics, SKIPPED_TABLE_FILES))
Assertions.assertEquals(resultedTableFiles, metric(metrics, RESULTED_TABLE_FILES))
Assertions.assertEquals(resultedRecordCount, metric(metrics, RESULTED_RECORD_COUNT))
Assertions.assertTrue(metric(metrics, RESULTED_TABLE_FILES_SIZE) > 0)
}

checkMetrics(s"SELECT * FROM T", 3, 0, 5)
checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 3, 2, 3)
checkMetrics(s"SELECT * FROM T", 3, 0, 5, 5)
checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 3, 2, 3, 3)

sql(s"DELETE FROM T WHERE pt = 'p1'")
checkMetrics(s"SELECT * FROM T", 4, 0, 4)
checkMetrics(s"SELECT * FROM T", 4, 0, 4, 4)

sql("CALL sys.compact(table => 'T', partitions => 'pt=\"p2\"')")
checkMetrics(s"SELECT * FROM T", 5, 0, 2)
checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 5, 1, 1)
checkMetrics(s"SELECT * FROM T", 5, 0, 2, 4)
checkMetrics(s"SELECT * FROM T WHERE pt = 'p2'", 5, 1, 1, 3)

// a scan without any filter cannot prune any manifest
val fullScan = getPaimonScan(s"SELECT * FROM T")
fullScan.inputPartitions
Assertions.assertEquals(0, metric(fullScan.reportDriverMetrics(), SKIPPED_MANIFESTS))
}
}

Expand Down
Loading