diff --git a/cdap-messaging-ext-spanner/src/main/java/io/cdap/cdap/messaging/spanner/SpannerMessagingService.java b/cdap-messaging-ext-spanner/src/main/java/io/cdap/cdap/messaging/spanner/SpannerMessagingService.java index 1dee6f1749df..f1c3edd51730 100644 --- a/cdap-messaging-ext-spanner/src/main/java/io/cdap/cdap/messaging/spanner/SpannerMessagingService.java +++ b/cdap-messaging-ext-spanner/src/main/java/io/cdap/cdap/messaging/spanner/SpannerMessagingService.java @@ -187,12 +187,15 @@ private String getCreateTopicMetadataDDLStatement() { *
*/ private String getCreateTopicDDLStatement(TopicId topicId) { - return String.format("CREATE TABLE IF NOT EXISTS %s ( %s INT64, %s INT64, %s" - + " TIMESTAMP NOT NULL OPTIONS (allow_commit_timestamp=true), %s INT64, %s BYTES(MAX) )" - + " PRIMARY KEY (%s, %s, %s), ROW DELETION POLICY" + " (OLDER_THAN(%s, INTERVAL 7 DAY))", - getTableName(topicId), SEQUENCE_ID_FIELD, PAYLOAD_SEQUENCE_ID_FIELD, PUBLISH_TS_FIELD, - PAYLOAD_REMAINING_CHUNKS_FIELD, PAYLOAD_FIELD, SEQUENCE_ID_FIELD, PAYLOAD_SEQUENCE_ID_FIELD, - PUBLISH_TS_FIELD, PUBLISH_TS_FIELD); + return String.format( + "CREATE TABLE IF NOT EXISTS %s ( %s TIMESTAMP NOT NULL OPTIONS (allow_commit_timestamp=true)," + + " %s INT64, %s INT64, %s INT64, %s BYTES(MAX) )" + + " PRIMARY KEY (%s, %s, %s), ROW DELETION POLICY" + + " (OLDER_THAN(%s, INTERVAL 7 DAY))", getTableName(topicId), + PUBLISH_TS_FIELD, SEQUENCE_ID_FIELD, PAYLOAD_SEQUENCE_ID_FIELD, + PAYLOAD_REMAINING_CHUNKS_FIELD, PAYLOAD_FIELD, + PUBLISH_TS_FIELD, SEQUENCE_ID_FIELD, PAYLOAD_SEQUENCE_ID_FIELD, + PUBLISH_TS_FIELD); } private void updateTopicMetadataTable(List