【问题标题】:Notifying postgres changes to java application将 postgres 更改通知给 java 应用程序
【发布时间】:2013-08-10 04:11:56
【问题描述】:

问题

我正在为几十万种产品构建 postgres 数据库。我将设置一个索引(Solr 或 ElasticSearch),以缩短复杂搜索查询的查询时间。

现在的重点是如何让索引与数据库同步?

过去我有一种应用程序会定期轮询数据库以检查应该完成的更新,但我会有一个过时的索引状态时间(从数据库更新到索引更新拉)。

我更喜欢这样的解决方案,在该解决方案中,数据库会通知我的应用程序(java 应用程序)数据库中的某些内容已发生更改,然后应用程序将决定是否需要更改索引更新与否。更准确地说,我会构建一种生产者和消费者结构,希望副本能够从 postgres 收到通知,告知发生了一些变化,如果这与索引的数据相关,它会存储在要更新的堆栈中。消费者将使用此堆栈并构建要存储到索引中的文档。

可能的解决方案

一种解决方案是编写一种副本端点,在该端点中,应用程序将充当 postgres 实例,用于从原始数据库复制数据。有人对这种方法有一些经验吗?

我还有哪些其他解决方案可以解决这个问题?

【问题讨论】:

    标签: java postgresql indexing synchronization replication


    【解决方案1】:

    对于这个问题我还有什么其他的解决方案?

    使用LISTEN and NOTIFY 告诉您的应用程序发生了变化。

    您可以从触发器发送NOTIFY,该触发器还记录队列表中的更改。

    您需要一个 PgJDBC 连接,该连接已为您正在使用的事件发送了 LISTEN。如果您使用 SSL,它必须通过定期发送空查询 ("") 来轮询数据库;如果您不使用 SSL,则可以通过使用异步通知检查来避免这种情况。您需要从连接池中解开 Connection 对象,以便能够将底层连接转换为 PgConnection 以使用监听/通知。见related answer

    生产者/消费者位会更难。要在 PostgreSQL 中拥有多个崩溃安全的并发使用者,您需要使用带有pg_try_advisory_lock(...) 的咨询锁定。如果您不需要并发消费者,那么这很容易,您只需 SELECT ... LIMIT 1 FOR UPDATE 一次一行。

    希望 9.4 将包含一个更简单的方法,使用 FOR UPDATE 跳过锁定的行,因为它正在开发中。

    【讨论】:

    • 只是为了记录 SKIP LOCKED 是在 PG 9.5 上引入的,避免使用像 pg_try_advisory_lock 这样的管理功能。
    • 是的,尽管咨询锁仍然有一些非常方便的用途。您不必再像上面那样在排队时使用它们,但它们在其他领域仍然非常方便。
    【解决方案2】:

    要使用 postgres 的 LISTEN 和 NOTIFY,您需要使用可以支持异步通知的驱动程序。 postgres JDBC 驱动不支持异步通知。

    要通过应用服务器的通道不断监听,请使用 pgjdbc-ng 0.6 驱动程序。

    http://impossibl.github.io/pgjdbc-ng/

    它支持异步通知,无需轮询。

    【讨论】:

      【解决方案3】:

      一般来说,我建议使用EAI patterns 来实现松散耦合。然后,如果你决定交换数据库,索引端的代码不会改变。

      如果你想坚持紧耦合,我建议使用 LISTEN/NOTIFY。 在 Java 中,使用 pgjdbc-ng driver 很重要,因为它支持异步 无需轮询的通知。

      这是一个异步模式(基于this answer):

      import com.impossibl.postgres.api.jdbc.PGConnection;
      import com.impossibl.postgres.api.jdbc.PGNotificationListener;
      import com.impossibl.postgres.jdbc.PGDataSource;    
      import java.sql.Statement;
      
      public static void listenToNotifyMessage() {
          PGDataSource dataSource = new PGDataSource();
          dataSource.setHost("localhost");
          dataSource.setPort(5432);
          dataSource.setDatabase("database_name");
          dataSource.setUser("postgres");
          dataSource.setPassword("password");
      
          PGNotificationListener listener = (int processId, String channelName, String payload) 
             -> System.out.println("notification = " + payload);
      
          try (PGConnection connection = (PGConnection) dataSource.getConnection()) {
              Statement statement = connection.createStatement();
              statement.execute("LISTEN test");
              statement.close();
              connection.addNotificationListener(listener);
              // it only works if the connection is open. Therefore, we do an endless loop here.
              while (true) {
                 Thread.sleep(500);
             }
          } catch (Exception e) {
              System.err.println(e);
          }
      }
      

      在其他语句中,您现在可以执行NOTIFY test, 'This is a payload';。您还可以在触发器等中执行NOTIFY

      【讨论】:

        猜你喜欢
        • 2013-03-28
        • 1970-01-01
        • 2016-07-09
        • 1970-01-01
        • 2021-01-18
        • 1970-01-01
        • 2020-08-20
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多