diff --git a/src/backend/commands/publicationcmds.c b/src/backend/commands/publicationcmds.c
index 301050104dd..9015a5a5136 100644
--- a/src/backend/commands/publicationcmds.c
+++ b/src/backend/commands/publicationcmds.c
@@ -940,6 +940,14 @@ CreatePublication(ParseState *pstate, CreatePublicationStmt *stmt)
 
 	if (stmt->for_all_tables)
 	{
+		/*
+		 * Block concurrent data-modifying statements that may rely on this
+		 * table not being covered by a FOR ALL TABLES publication requiring a
+		 * replica identity. This conflicts with the RowExclusiveLock taken by
+		 * CheckCmdReplicaIdentity() on the same object.
+ 		 */
+		LockRelationOid(PublicationRelationId, ShareRowExclusiveLock);
+
 		/* Process EXCEPT table list */
 		if (exceptrelations != NIL)
 		{
@@ -1610,6 +1618,15 @@ AlterPublicationAllFlags(AlterPublicationStmt *stmt, Relation rel,
 	/* Update FOR ALL TABLES flag if changed */
 	if (stmt->for_all_tables != pubform->puballtables)
 	{
+		/*
+		 * Block concurrent data-modifying statements when setting FOR ALL
+		 * TABLES publication requiring a replica identity. This conflicts with
+		 * the RowExclusiveLock taken by CheckCmdReplicaIdentity() on the same
+		 * object.
+		 */
+		if (stmt->for_all_tables)
+			LockRelationOid(PublicationRelationId, ShareRowExclusiveLock);
+
 		values[Anum_pg_publication_puballtables - 1] =
 			BoolGetDatum(stmt->for_all_tables);
 		replaces[Anum_pg_publication_puballtables - 1] = true;
@@ -2005,8 +2022,9 @@ CloseTableList(List *rels)
 }
 
 /*
- * Lock the schemas specified in the schema list in AccessShareLock mode in
- * order to prevent concurrent schema deletion.
+ * Lock the schemas in ShareRowExclusiveLock mode to prevent concurrent
+ * schema deletion and data-modifying statements from racing with the
+ * publication change.
  */
 static void
 LockSchemaList(List *schemalist)
@@ -2019,7 +2037,7 @@ LockSchemaList(List *schemalist)
 
 		/* Allow query cancel in case this takes a long time */
 		CHECK_FOR_INTERRUPTS();
-		LockDatabaseObject(NamespaceRelationId, schemaid, 0, AccessShareLock);
+		LockDatabaseObject(NamespaceRelationId, schemaid, 0, ShareRowExclusiveLock);
 
 		/*
 		 * It is possible that by the time we acquire the lock on schema,
diff --git a/src/backend/executor/execReplication.c b/src/backend/executor/execReplication.c
index dd42acc13e2..fe7b229e908 100644
--- a/src/backend/executor/execReplication.c
+++ b/src/backend/executor/execReplication.c
@@ -24,6 +24,7 @@
 #include "access/xact.h"
 #include "access/heapam.h"
 #include "catalog/pg_am_d.h"
+#include "catalog/pg_namespace.h"
 #include "commands/trigger.h"
 #include "executor/executor.h"
 #include "executor/nodeModifyTable.h"
@@ -1031,6 +1032,7 @@ void
 CheckCmdReplicaIdentity(Relation rel, CmdType cmd)
 {
 	PublicationDesc pubdesc;
+	bool		has_local_identity;
 
 	/*
 	 * Skip checking the replica identity for partitioned tables, because the
@@ -1043,6 +1045,30 @@ CheckCmdReplicaIdentity(Relation rel, CmdType cmd)
 	if (cmd != CMD_UPDATE && cmd != CMD_DELETE)
 		return;
 
+	/* If relation has replica identity we are always good. */
+	has_local_identity = OidIsValid(RelationGetReplicaIndex(rel)) ||
+		rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL;
+
+	/*
+	 * Tables without a local replica identity can race with publication DDL
+	 * for FOR ALL TABLES and TABLES IN SCHEMA. Take locks on the table's
+	 * schema and the publication catalog before using the publication
+	 * descriptor.
+	 *
+	 * The locks conflict with the corresponding publication DDL locks and also
+	 * ensure that any pending catalog invalidations are processed before we
+	 * build the descriptor.
+	 *
+	 * Tables with a local replica identity are not affected by this race, so
+	 * they do not need these additional locks.
+	 */
+	if (!has_local_identity)
+	{
+		LockDatabaseObject(NamespaceRelationId, RelationGetNamespace(rel), 0,
+						   RowExclusiveLock);
+		LockRelationOid(PublicationRelationId, RowExclusiveLock);
+	}
+
 	/*
 	 * It is only safe to execute UPDATE/DELETE if the relation does not
 	 * publish UPDATEs or DELETEs, or all the following conditions are
@@ -1104,12 +1130,7 @@ CheckCmdReplicaIdentity(Relation rel, CmdType cmd)
 						RelationGetRelationName(rel)),
 				 errdetail("Replica identity must not contain unpublished generated columns.")));
 
-	/* If relation has replica identity we are always good. */
-	if (OidIsValid(RelationGetReplicaIndex(rel)))
-		return;
-
-	/* REPLICA IDENTITY FULL is also good for UPDATE/DELETE. */
-	if (rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL)
+	if (has_local_identity)
 		return;
 
 	/*
