4646ROW_WATERMARK = "row"
4747WATERMARK_PREFIXES = (TOMBSTONE_WATERMARK , ROW_WATERMARK )
4848
49+ WRITE_WATERMARK_TO_POSTGRES_OPTION = "hybrid_cloud.write_deletion_watermark_to_postgres"
50+ READ_WATERMARK_FROM_POSTGRES_OPTION = "hybrid_cloud.read_deletion_watermark_from_postgres"
51+
4952
5053@dataclass
5154class WatermarkBatch :
@@ -96,7 +99,7 @@ def _write_watermark(
9699 )
97100
98101 # Dual-write deletion watermarks to Redis and Postgres
99- if options .get ("hybrid_cloud.write_deletion_watermark_to_postgres" ):
102+ if options .get (WRITE_WATERMARK_TO_POSTGRES_OPTION ):
100103 try :
101104 _watermark_model (field ).objects .update_or_create (
102105 prefix = prefix ,
@@ -116,20 +119,54 @@ def _write_watermark(
116119 )
117120
118121
119- def get_watermark (prefix : str , field : HybridCloudForeignKey [Any , Any ]) -> tuple [int , str ]:
120- client = _get_redis_client ()
121- key = get_watermark_key (prefix , field )
122- v = client .get (key )
122+ def _watermark_row_lookup (prefix : str , field : HybridCloudForeignKey [Any , Any ]) -> dict [str , str ]:
123+ return dict (
124+ prefix = prefix ,
125+ table_name = field .model ._meta .db_table ,
126+ field_name = field .name ,
127+ )
128+
129+
130+ def _read_redis_watermark (
131+ prefix : str , field : HybridCloudForeignKey [Any , Any ]
132+ ) -> tuple [int , str ] | None :
133+ v = _get_redis_client ().get (get_watermark_key (prefix , field ))
123134 if v is None :
124- result = (0 , uuid4 ().hex )
125- _write_watermark (prefix , field , * result )
126- return result
135+ return None
127136 lower , transaction_id = json .loads (v )
128137 if not (isinstance (lower , int ) and isinstance (transaction_id , str )):
129138 raise TypeError ("Expected watermarks data to be a tuple of (int, str)" )
130139 return lower , transaction_id
131140
132141
142+ def _read_postgres_watermark (
143+ prefix : str , field : HybridCloudForeignKey [Any , Any ]
144+ ) -> tuple [int , str ] | None :
145+ row = _watermark_model (field ).objects .filter (** _watermark_row_lookup (prefix , field )).first ()
146+ if row is None :
147+ return None
148+ return row .low_bound , row .transaction_id
149+
150+
151+ def _postgres_is_source_of_truth () -> bool :
152+ return bool (options .get (WRITE_WATERMARK_TO_POSTGRES_OPTION )) and bool (
153+ options .get (READ_WATERMARK_FROM_POSTGRES_OPTION )
154+ )
155+
156+
157+ def get_watermark (prefix : str , field : HybridCloudForeignKey [Any , Any ]) -> tuple [int , str ]:
158+ if _postgres_is_source_of_truth ():
159+ watermark = _read_postgres_watermark (prefix , field )
160+ else :
161+ watermark = _read_redis_watermark (prefix , field )
162+
163+ if watermark is None :
164+ result = (0 , uuid4 ().hex )
165+ _write_watermark (prefix , field , * result )
166+ return result
167+ return watermark
168+
169+
133170def set_watermark (
134171 prefix : str , field : HybridCloudForeignKey [Any , Any ], value : int , prev_transaction_id : str
135172) -> None :
0 commit comments