On Tue, Aug 4, 2026 at 1:05 AM Chao Li <[email protected]> wrote:
>
>
>
> > On Aug 4, 2026, at 08:41, Masahiko Sawada <[email protected]> wrote:
> >
> > On Thu, Jul 9, 2026 at 5:14 AM Fujii Masao <[email protected]> wrote:
> >>
> >> On Wed, Jun 24, 2026 at 3:02 AM Masahiko Sawada <[email protected]> 
> >> wrote:
> >>> Right. I'm implementing a basic DDL replication solution using this
> >>> hook as a sample implementation. Other than that, this feature can be
> >>> used to add additional information for the replicated transaction that
> >>> is not replicated via logical replication protocol, such as the
> >>> executed user or context-specific information, which can then be
> >>> dispatched to a CDC system connected to the subscriber.
> >>
> >> +1 for the proposed hook, although I haven't read the patch yet.
> >>
> >> I know of a system that uses logical decoding messages to propagate data
> >> to external systems. In that system, the application writes the data to
> >> be propagated as logical decoding messages, and a CDC pipeline consumes
> >> them via logical decoding and sends them to other systems, for example
> >> through Kafka.
> >>
> >> In that system, when the remote site is maintained by physical replication,
> >> the same logical decoding messages can be decoded on the standby
> >> at the remote site to deliver them to external systems there. However,
> >> if the remote site uses logical replication instead, those messages are
> >> currently neither delivered to nor processed on the subscriber at
> >> the remote site. As a result, the CDC pipeline at the remote site cannot
> >> consume the data they carry.
> >>
> >> The proposed hook might be useful for this use case. An extension could
> >> process the incoming logical decoding messages on the subscriber and
> >> forward them to the local CDC pipeline, or store or re-emit them in a form
> >> that local consumers can process.
> >
> > Thank you for sharing the concrete use case. This is one of the use
> > cases I initially imagined.
> >
> > I've rebased and updated the patch. Please review it.
> >
> > Regards,
> >
> > --
> > Masahiko Sawada
> > Amazon Web Services: https://aws.amazon.com
> > <v2-0001-Add-a-hook-for-handling-logical-messages-on-subsc.patch>
>
> I played with this feature today. I have a few comments from a design 
> perspective:

Thank you for reviewing the patch!

> 1. Would it make sense to pass the receiving subscription’s identity 
> explicitly to the hook? Otherwise, a hook that needs to distinguish messages 
> received by different subscriptions has to obtain it from a global variable 
> such as MySubscription->oid.

It's common for hook functions to access a global variable to get more
information such as MyDatabaseId and MyProcPort etc. Using
MySubscription to get a subscription's identity works for me.

> 2. A hook is invoked by an apply worker, and since a message has no target 
> relation, the hook runs as the subscription owner regardless of run_as_owner. 
> I think it would be useful to document this explicitly.

I think it depends on the hook function implementations. They can
switch the role as they want.

>
> 3. I’m thinking out loud here. Would it be useful to provide an opt-in 
> default hook that re-emits received messages, so that they are written to the 
> subscriber-side WAL and can be consumed by a local logical-decoding client? I 
> understand that a re-emitted message would have a different LSN, just as 
> logically replicated row changes generate new local WAL records. For the use 
> case Fujii-san described, this might allow an existing decoder to continue 
> working after the remote site switches from physical to logical replication.

Interesting idea, but I'm not sure we should have it in the core or in
the contrib. I think it's a good topic to discuss in a separate
thread.

>
> And a few comment for the code changes:
>
> 1 - worker_internal.h
> ```
> +typedef void (*LogicalRepMessageHandle_hook_type) (LogicalRepMessageData 
> *msg);
> +extern PGDLLIMPORT LogicalRepMessageHandle_hook_type 
> LogicalRepMessageHandle_hook;
> ```
>
> I guess we don’t expect a hook function to mutate the message, so maybe make 
> the pointer const.
>
> 2 - worker.c
> ```
> +       /*
> +        * Logical messages are relation-agnostic, so they don't map cleanly 
> onto
> +        * the tablesync worker of a particular relation, and there would be 
> no
> +        * well-defined ordering between a message and the initial copy. Leave
> +        * them to the (parallel) apply workers. Logical messages are handled 
> only
> +        * the (parallel) apply workers
> +        */
> ```
>
> The last sentence looks duplicate and incomplete.
>
> 3 - proto.c
> ```
> +       /* read message length */
> +       len = pq_getmsgint(in, 4);
> +       msg_data->message_size = len;
> +
> +       /* and data */
> +       msg = palloc(len + 1);
> +       pq_copymsgbytes(in, msg, len);
> +
> +       msg[len] = '\0';
> +       msg_data->message = msg;
> ```
>
> logicalrep_read_message() allocates one extra byte for NULL terminator, which 
> is unnecessary, as the contract is to use msg_data->message_size to decide 
> the message length, and a message contain contains 0 in the middle. But I 
> agree it may be useful for debugging and logging, so maybe add a comment to 
> explain why using this extra byte.

Agreed with the all above comments.

I'll submit the updated patch shortly.

Regards,

-- 
Masahiko Sawada
Amazon Web Services: https://aws.amazon.com


Reply via email to