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 @@ -92,7 +92,7 @@ public abstract class MetaDataProtocol extends MetaDataService {
public static final long MIN_SYSTEM_TABLE_TIMESTAMP_5_1_0 = MIN_SYSTEM_TABLE_TIMESTAMP_4_16_0;
public static final long MIN_SYSTEM_TABLE_TIMESTAMP_5_2_0 = MIN_TABLE_TIMESTAMP + 38;
public static final long MIN_SYSTEM_TABLE_TIMESTAMP_5_3_0 = MIN_TABLE_TIMESTAMP + 44;
public static final long MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 = MIN_TABLE_TIMESTAMP + 45;
public static final long MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 = MIN_TABLE_TIMESTAMP + 46;
// MIN_SYSTEM_TABLE_TIMESTAMP needs to be set to the max of all the MIN_SYSTEM_TABLE_TIMESTAMP_*
// constants
public static final long MIN_SYSTEM_TABLE_TIMESTAMP = MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,18 @@ public class PhoenixDatabaseMetaData implements DatabaseMetaData {
public static final byte[] INDEX_TYPE_BYTES = Bytes.toBytes(INDEX_TYPE);
public static final String INDEX_CONSISTENCY = "INDEX_CONSISTENCY";
public static final byte[] INDEX_CONSISTENCY_BYTES = Bytes.toBytes(INDEX_CONSISTENCY);
// No-op SYSTEM.CATALOG marker column, never populated or read. It exists solely to advance the
// SYSTEM.CATALOG header timestamp to MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 during an in-place upgrade.
// The isUpgradeRequired() gate keys off SYSTEM.CATALOG's header timestamp; the paired
// SYSTEM.TRANSFORM column-add does not itself advance that timestamp, so on a cluster
// bootstrapped
// in the window between sibling 5.4.0 features (already at the old threshold via
// INDEX_CONSISTENCY)
// the transform column-add would be unreachable and clients would loop on
// UpgradeRequiredException.
// Adding this genuinely-new column at the new threshold makes the upgrade reachable. Single-use:
// a future timestamp bump needs its own new marker (re-adding an existing column is a no-op).
public static final String UPGRADE_TS_ANCHOR_5_4_0 = "UPGRADE_TS_ANCHOR_5_4_0";
public static final String LINK_TYPE = "LINK_TYPE";
public static final byte[] LINK_TYPE_BYTES = Bytes.toBytes(LINK_TYPE);
public static final String TASK_TYPE = "TASK_TYPE";
Expand All @@ -236,6 +248,14 @@ public class PhoenixDatabaseMetaData implements DatabaseMetaData {
public static final String OLD_METADATA = "OLD_METADATA";
public static final String NEW_METADATA = "NEW_METADATA";
public static final String TRANSFORM_FUNCTION = "TRANSFORM_FUNCTION";
// Epoch-millis (BIGINT, nullable) marking the earliest time the transform monitor may leave the
// PENDING_PARTIAL_PASS wait window. Compared against EnvironmentEdgeManager.currentTimeMillis().
public static final String PENDING_PARTIAL_PASS_UNTIL_TS = "PENDING_PARTIAL_PASS_UNTIL_TS";
// Epoch-millis (BIGINT, nullable) captured just before the cutover pointer swap and preserved
// across later transitions. The partial pass derives its repair-scan lower bound from this so
// rows written to the old pointer during the post-cutover cache-refresh window are re-verified
// rather than skipped.
public static final String CUTOVER_TS = "CUTOVER_TS";
public static final String TRANSFORM_TABLE_TTL = "7776000"; // 90 days

public static final int TTL_FOR_MUTEX = 15 * 60; // 15min
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4801,39 +4801,46 @@ protected PhoenixConnection upgradeSystemCatalogIfRequired(PhoenixConnection met
}
if (currentServerSideTableTimeStamp < MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0) {
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 9,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 10,
PhoenixDatabaseMetaData.PHYSICAL_TABLE_NAME + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 8,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 9,
PhoenixDatabaseMetaData.SCHEMA_VERSION + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 7,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 8,
PhoenixDatabaseMetaData.EXTERNAL_SCHEMA_ID + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 6,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 7,
PhoenixDatabaseMetaData.STREAMING_TOPIC_NAME + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 5,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 6,
PhoenixDatabaseMetaData.INDEX_WHERE + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 4,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 5,
PhoenixDatabaseMetaData.CDC_INCLUDE_TABLE + " " + PVarchar.INSTANCE.getSqlTypeName());

/**
* TODO: Provide a path to copy existing data from PHOENIX_TTL to TTL column and then to DROP
* PHOENIX_TTL Column. See PHOENIX-7023
*/
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 3,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 4,
PhoenixDatabaseMetaData.TTL + " " + PVarchar.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 2,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 3,
PhoenixDatabaseMetaData.ROW_KEY_MATCHER + " " + PVarbinary.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 1,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 2,
PhoenixDatabaseMetaData.IS_STRICT_TTL + " " + PBoolean.INSTANCE.getSqlTypeName());
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0, PhoenixDatabaseMetaData.INDEX_CONSISTENCY + " CHAR(1)");
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 - 1,
PhoenixDatabaseMetaData.INDEX_CONSISTENCY + " CHAR(1)");
// No-op catalog schema-version anchor: advances the SYSTEM.CATALOG header timestamp to the
// new MIN so in-place SNAPSHOT clusters re-enter the upgrade path and pick up the new
// SYSTEM.TRANSFORM columns. Never populated or read.
metaConnection = addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_CATALOG,
MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0,
PhoenixDatabaseMetaData.UPGRADE_TS_ANCHOR_5_4_0 + " CHAR(1)");

// move TTL values stored in descriptor to SYSCAT TTL column.
moveTTLFromHBaseLevelTTLToPhoenixLevelTTL(metaConnection);
Expand Down Expand Up @@ -5297,12 +5304,28 @@ private PhoenixConnection upgradeSystemTask(PhoenixConnection metaConnection,
return metaConnection;
}

private PhoenixConnection upgradeSystemTransform(PhoenixConnection metaConnection,
@VisibleForTesting
public PhoenixConnection upgradeSystemTransform(PhoenixConnection metaConnection,
Map<String, String> systemTableToSnapshotMap) throws SQLException {
try (Statement statement = metaConnection.createStatement()) {
statement.executeUpdate(getTransformDDL());
} catch (TableAlreadyExistsException ignored) {

} catch (NewerTableAlreadyExistsException ignored) {
// A newer SYSTEM.TRANSFORM header means a same-or-newer client already ran this DDL, whose
// CREATE statement carries the two new columns; the column-add below is therefore already
// done and skipping it is safe.
} catch (TableAlreadyExistsException e) {
// This is the first-ever column add to SYSTEM.TRANSFORM, so take a snapshot before altering.
takeSnapshotOfSysTable(systemTableToSnapshotMap, e);
// addColumnsIfNotExists is idempotent, so call it unconditionally rather than gating on the
// table timestamp. A gate keyed on MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0 is unreachable on
// SNAPSHOT clusters whose SYSTEM.TRANSFORM header already reached that timestamp without the
// columns, and would strand transform reads on a missing-column error.
metaConnection =
addColumnsIfNotExists(metaConnection, PhoenixDatabaseMetaData.SYSTEM_TRANSFORM_NAME,
MetaDataProtocol.MIN_SYSTEM_TABLE_TIMESTAMP_5_4_0,
PhoenixDatabaseMetaData.PENDING_PARTIAL_PASS_UNTIL_TS + " "
+ PLong.INSTANCE.getSqlTypeName() + ", " + PhoenixDatabaseMetaData.CUTOVER_TS + " "
+ PLong.INSTANCE.getSqlTypeName());
}
return metaConnection;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.COLUMN_QUALIFIER_COUNTER;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.COLUMN_SIZE;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.CURRENT_VALUE;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.CUTOVER_TS;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.CYCLE_FLAG;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.DATA_TABLE_NAME;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.DATA_TYPE;
Expand Down Expand Up @@ -99,6 +100,7 @@
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PARTITION_ID;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PARTITION_START_KEY;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PARTITION_START_TIME;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PENDING_PARTIAL_PASS_UNTIL_TS;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PHOENIX_TTL;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PHOENIX_TTL_HWM;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.PHYSICAL_NAME;
Expand Down Expand Up @@ -172,6 +174,7 @@
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.TYPE_NAME;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.TYPE_SEQUENCE;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.UPDATE_CACHE_FREQUENCY;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.UPGRADE_TS_ANCHOR_5_4_0;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.USER;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.USE_STATS_FOR_PARALLELIZATION;
import static org.apache.phoenix.jdbc.PhoenixDatabaseMetaData.VIEW_CONSTANT;
Expand Down Expand Up @@ -419,7 +422,8 @@ enum JoinType {
+ " BOOLEAN, \n" + SCHEMA_VERSION + " VARCHAR, \n" + EXTERNAL_SCHEMA_ID + " VARCHAR, \n"
+ STREAMING_TOPIC_NAME + " VARCHAR, \n" + INDEX_WHERE + " VARCHAR, \n" + CDC_INCLUDE_TABLE
+ " VARCHAR, \n" + TTL + " VARCHAR, \n" + ROW_KEY_MATCHER + " VARBINARY_ENCODED, \n"
+ IS_STRICT_TTL + " BOOLEAN, \n" + INDEX_CONSISTENCY + " CHAR(1), \n" +
+ IS_STRICT_TTL + " BOOLEAN, \n" + INDEX_CONSISTENCY + " CHAR(1), \n"
+ UPGRADE_TS_ANCHOR_5_4_0 + " CHAR(1), \n" +
// Column metadata (will be null for table row)
DATA_TYPE + " INTEGER," + COLUMN_SIZE + " INTEGER," + DECIMAL_DIGITS + " INTEGER," + NULLABLE
+ " INTEGER," + ORDINAL_POSITION + " INTEGER," + SORT_ORDER + " INTEGER," + ARRAY_SIZE
Expand Down Expand Up @@ -566,9 +570,10 @@ enum JoinType {
TRANSFORM_STATUS + " VARCHAR NULL," + TRANSFORM_JOB_ID + " VARCHAR NULL,"
+ TRANSFORM_RETRY_COUNT + " INTEGER NULL," + TRANSFORM_START_TS + " TIMESTAMP NULL,"
+ TRANSFORM_LAST_STATE_TS + " TIMESTAMP NULL," + OLD_METADATA + " VARBINARY NULL,\n"
+ NEW_METADATA + " VARCHAR NULL,\n" + TRANSFORM_FUNCTION + " VARCHAR NULL\n" + "CONSTRAINT "
+ SYSTEM_TABLE_PK_NAME + " PRIMARY KEY (" + TENANT_ID + "," + TABLE_SCHEM + ","
+ LOGICAL_TABLE_NAME + "))\n" + HConstants.VERSIONS + "=%s,\n"
+ NEW_METADATA + " VARCHAR NULL,\n" + TRANSFORM_FUNCTION + " VARCHAR NULL,\n"
+ PENDING_PARTIAL_PASS_UNTIL_TS + " BIGINT NULL,\n" + CUTOVER_TS + " BIGINT NULL\n"
+ "CONSTRAINT " + SYSTEM_TABLE_PK_NAME + " PRIMARY KEY (" + TENANT_ID + "," + TABLE_SCHEM
+ "," + LOGICAL_TABLE_NAME + "))\n" + HConstants.VERSIONS + "=%s,\n"
+ ColumnFamilyDescriptorBuilder.KEEP_DELETED_CELLS + "=%s,\n"
+ ColumnFamilyDescriptorBuilder.TTL + "=" + TRANSFORM_TABLE_TTL + ",\n" + // 90 days
TableDescriptorBuilder.SPLIT_POLICY + "='" + SYSTEM_TASK_SPLIT_POLICY_CLASSNAME + "',\n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -297,7 +297,11 @@ public static TransformType fromSerializedValue(int serializedValue) {
}

public static TransformType getPartialTransform(TransformType transformType) {
if (transformType == METADATA_TRANSFORM) {
// A full transform's partial variant is the partial type. Asking for the partial variant of a
// record that is already partial-type yields the partial type itself: this is what lets a
// failed partial pass be re-kicked (the record is already partial-type on retry) instead of
// being mistaken for "no partial pass needed" and completed early.
if (transformType == METADATA_TRANSFORM || transformType == METADATA_TRANSFORM_PARTIAL) {
return METADATA_TRANSFORM_PARTIAL;
}
return null;
Expand Down Expand Up @@ -326,6 +330,20 @@ public String toString() {
return "PENDING_CUTOVER";
}
},
/**
* Cutover is done; the post-cutover partial-pass repair scan is deferred until the UCF wait.
*/
PENDING_PARTIAL_PASS {
public String toString() {
return "PENDING_PARTIAL_PASS";
}
},
/** The post-cutover partial-pass repair scan has been launched and is being monitored. */
PARTIAL_PASS_RUNNING {
public String toString() {
return "PARTIAL_PASS_RUNNING";
}
},
COMPLETED {
public String toString() {
return "COMPLETED";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,11 +45,14 @@ public class SystemTransformRecord {
private final byte[] oldMetadata;
private final String newMetadata;
private final String transformFunction;
private final Long pendingPartialPassUntilTs;
private final Long cutoverTs;

public SystemTransformRecord(PTable.TransformType transformType, String schemaName,
String logicalTableName, String tenantId, String newPhysicalTableName, String logicalParentName,
String transformStatus, String transformJobId, Integer transformRetryCount, Timestamp startTs,
Timestamp lastStateTs, byte[] oldMetadata, String newMetadata, String transformFunction) {
Timestamp lastStateTs, byte[] oldMetadata, String newMetadata, String transformFunction,
Long pendingPartialPassUntilTs, Long cutoverTs) {
this.transformType = transformType;
this.schemaName = schemaName;
this.tenantId = tenantId;
Expand All @@ -64,14 +67,17 @@ public SystemTransformRecord(PTable.TransformType transformType, String schemaNa
this.oldMetadata = oldMetadata;
this.newMetadata = newMetadata;
this.transformFunction = transformFunction;
this.pendingPartialPassUntilTs = pendingPartialPassUntilTs;
this.cutoverTs = cutoverTs;
}

public String getString() {
return String.format(
"transformType: %s, schameName: %s, logicalTableName: %s, newPhysicalTableName: %s, logicalParentName: %s, status: %s",
"transformType: %s, schameName: %s, logicalTableName: %s, newPhysicalTableName: %s, logicalParentName: %s, status: %s, pendingPartialPassUntilTs: %s, cutoverTs: %s",
String.valueOf(transformType), String.valueOf(schemaName), String.valueOf(logicalTableName),
String.valueOf(newPhysicalTableName), String.valueOf(logicalParentName),
String.valueOf(transformStatus));
String.valueOf(transformStatus), String.valueOf(pendingPartialPassUntilTs),
String.valueOf(cutoverTs));
}

public PTable.TransformType getTransformType() {
Expand Down Expand Up @@ -130,10 +136,20 @@ public String getTransformFunction() {
return transformFunction;
}

public Long getPendingPartialPassUntilTs() {
return pendingPartialPassUntilTs;
}

public Long getCutoverTs() {
return cutoverTs;
}

public boolean isActive() {
return (transformStatus.equals(PTable.TransformStatus.STARTED.name())
|| transformStatus.equals(PTable.TransformStatus.CREATED.name())
|| transformStatus.equals(PTable.TransformStatus.PENDING_CUTOVER.name()));
|| transformStatus.equals(PTable.TransformStatus.PENDING_CUTOVER.name())
|| transformStatus.equals(PTable.TransformStatus.PENDING_PARTIAL_PASS.name())
|| transformStatus.equals(PTable.TransformStatus.PARTIAL_PASS_RUNNING.name()));
}

@edu.umd.cs.findbugs.annotations.SuppressWarnings(value = { "EI_EXPOSE_REP", "EI_EXPOSE_REP2" },
Expand All @@ -154,6 +170,8 @@ public static class SystemTransformBuilder {
private byte[] oldMetadata;
private String newMetadata;
private String transformFunction;
private Long pendingPartialPassUntilTs;
private Long cutoverTs;

public SystemTransformBuilder() {

Expand All @@ -174,6 +192,8 @@ public SystemTransformBuilder(SystemTransformRecord systemTransformRecord) {
this.setOldMetadata(systemTransformRecord.getOldMetadata());
this.setNewMetadata(systemTransformRecord.getNewMetadata());
this.setTransformFunction(systemTransformRecord.getTransformFunction());
this.setPendingPartialPassUntilTs(systemTransformRecord.getPendingPartialPassUntilTs());
this.setCutoverTs(systemTransformRecord.getCutoverTs());
}

public SystemTransformBuilder setTransformType(PTable.TransformType transformType) {
Expand Down Expand Up @@ -246,6 +266,16 @@ public SystemTransformBuilder setTransformFunction(String transformFunction) {
return this;
}

public SystemTransformBuilder setPendingPartialPassUntilTs(Long pendingPartialPassUntilTs) {
this.pendingPartialPassUntilTs = pendingPartialPassUntilTs;
return this;
}

public SystemTransformBuilder setCutoverTs(Long cutoverTs) {
this.cutoverTs = cutoverTs;
return this;
}

public SystemTransformRecord build() {
Timestamp lastTs = lastStateTs;
if (
Expand All @@ -256,7 +286,8 @@ public SystemTransformRecord build() {
}
return new SystemTransformRecord(transformType, schemaName, logicalTableName, tenantId,
newPhysicalTableName, logicalParentName, transformStatus, transformJobId,
transformRetryCount, startTs, lastTs, oldMetadata, newMetadata, transformFunction);
transformRetryCount, startTs, lastTs, oldMetadata, newMetadata, transformFunction,
pendingPartialPassUntilTs, cutoverTs);
}

public static SystemTransformRecord build(ResultSet resultSet) throws SQLException {
Expand All @@ -276,6 +307,10 @@ public static SystemTransformRecord build(ResultSet resultSet) throws SQLExcepti
builder.setOldMetadata(resultSet.getBytes(col++));
builder.setNewMetadata(resultSet.getString(col++));
builder.setTransformFunction(resultSet.getString(col++));
long pendingPartialPassUntilTs = resultSet.getLong(col++);
builder.setPendingPartialPassUntilTs(resultSet.wasNull() ? null : pendingPartialPassUntilTs);
long cutoverTs = resultSet.getLong(col++);
builder.setCutoverTs(resultSet.wasNull() ? null : cutoverTs);

return builder.build();
}
Expand Down
Loading