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
Original file line number Diff line number Diff line change
Expand Up @@ -2681,9 +2681,22 @@ public class ConfigOptions {
+ "so that clients begin throttling before the storage engine is forced to throttle itself. "
+ "The gap between this value and the storage trigger forms the throttle ramp-up window.");

// ------------------------------------------------------------------------
// ConfigOptions for KV lazy open

public static final ConfigOption<Boolean> KV_LAZY_OPEN_ENABLED =
key("kv.lazy-open.enabled")
.booleanType()
.defaultValue(false)
.withDescription("Whether to enable KvTablet lazy open.");

public static final ConfigOption<Duration> KV_LAZY_OPEN_IDLE_TIMEOUT =
key("kv.lazy-open.idle-timeout")
.durationType()
.defaultValue(Duration.ofHours(24))
.withDescription(
"Idle time before an open KvTablet is eligible for release back to lazy state.");

// ConfigOptions for metrics
// ------------------------------------------------------------------------
public static final ConfigOption<List<String>> METRICS_REPORTERS =
key("metrics.reporters")
.stringType()
Expand All @@ -2704,9 +2717,7 @@ public class ConfigOptions {
+ "the CoordinatorServer) it is advisable to use a port range "
+ "like 9250-9260.");

// ------------------------------------------------------------------------
// ConfigOptions for prometheus push gateway reporter
// ------------------------------------------------------------------------
public static final ConfigOption<String> METRICS_REPORTER_PROMETHEUS_PUSHGATEWAY_HOST_URL =
key("metrics.reporter.prometheus-push.host-url")
.stringType()
Expand Down Expand Up @@ -2771,9 +2782,7 @@ public class ConfigOptions {
+ "The value is automatically redacted when the configuration "
+ "is logged or displayed.");

// ------------------------------------------------------------------------
// ConfigOptions for jmx reporter
// ------------------------------------------------------------------------
public static final ConfigOption<String> METRICS_REPORTER_JMX_HOST =
key("metrics.reporter.jmx.port")
.stringType()
Expand All @@ -2786,9 +2795,7 @@ public class ConfigOptions {
+ "the CoordinatorServer) it is advisable to use a port range "
+ "like 9990-9999.");

// ------------------------------------------------------------------------
// ConfigOptions for influxdb reporter
// ------------------------------------------------------------------------
public static final ConfigOption<String> METRICS_REPORTER_INFLUXDB_VERSION =
key("metrics.reporter.influxdb.version")
.stringType()
Expand Down Expand Up @@ -2830,9 +2837,7 @@ public class ConfigOptions {
.defaultValue(Duration.ofSeconds(10))
.withDescription("The interval of reporting metrics to InfluxDB.");

// ------------------------------------------------------------------------
// ConfigOptions for lakehouse storage
// ------------------------------------------------------------------------
public static final ConfigOption<Boolean> DATALAKE_ENABLED =
key("datalake.enabled")
.booleanType()
Expand All @@ -2853,9 +2858,7 @@ public class ConfigOptions {
"The datalake format used by Fluss as lakehouse storage. Currently, supported formats are Paimon, Iceberg, Hudi, and Lance. "
+ "In the future, more kinds of data lake format will be supported, such as DeltaLake.");

// ------------------------------------------------------------------------
// ConfigOptions for tiering service
// ------------------------------------------------------------------------

public static final ConfigOption<Boolean> LAKE_TIERING_AUTO_EXPIRE_SNAPSHOT =
key("lake.tiering.auto-expire-snapshot")
Expand All @@ -2877,9 +2880,7 @@ public class ConfigOptions {
+ "If not configured and the tiering service runs in a Flink job, Fluss uses "
+ "Flink's IO temporary directories with a 'fluss' child directory.");

// ------------------------------------------------------------------------
// ConfigOptions for fluss kafka
// ------------------------------------------------------------------------
public static final ConfigOption<Boolean> KAFKA_ENABLED =
key("kafka.enabled")
.booleanType()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,9 @@ public class MetricNames {
"preWriteBufferTruncateAsDuplicatedPerSecond";
public static final String KV_PRE_WRITE_BUFFER_TRUNCATE_AS_ERROR_RATE =
"preWriteBufferTruncateAsErrorPerSecond";
public static final String KV_TABLET_OPEN_COUNT = "kvTabletOpenCount";
public static final String KV_TABLET_LAZY_COUNT = "kvTabletLazyCount";
public static final String KV_TABLET_FAILED_COUNT = "kvTabletFailedCount";

// --------------------------------------------------------------------------------------------
// RocksDB metrics
Expand Down
13 changes: 1 addition & 12 deletions fluss-filesystems/fluss-fs-oss/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,7 @@
<execution>
<goals>
<goal>jar</goal>
<goal>test-jar</goal>
</goals>
</execution>
</executions>
Expand Down Expand Up @@ -237,18 +238,6 @@
</execution>
</executions>
</plugin>

<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-jar-plugin</artifactId>
<executions>
<execution>
<goals>
<goal>test-jar</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>

Expand Down
7 changes: 0 additions & 7 deletions fluss-flink/fluss-flink-common/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -117,13 +117,6 @@
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
<version>${flink.minor.version}</version>
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-base</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,10 @@ void testBoundedPkTableEmitsKvBatchSplits() throws Throwable {

@Test
void testBoundedPkTableEmitsSnapshotSplitsByDefault() throws Throwable {
createTable(DEFAULT_TABLE_PATH, DEFAULT_PK_TABLE_DESCRIPTOR);
long tableId = createTable(DEFAULT_TABLE_PATH, DEFAULT_PK_TABLE_DESCRIPTOR);
for (int bucketId = 0; bucketId < DEFAULT_BUCKET_NUM; bucketId++) {
FLUSS_CLUSTER_EXTENSION.waitUntilAllReplicaReady(new TableBucket(tableId, bucketId));
}
int numSubtasks = DEFAULT_BUCKET_NUM;
try (MockSplitEnumeratorContext<SourceSplitBase> context =
new MockSplitEnumeratorContext<>(numSubtasks)) {
Expand Down
Loading
Loading