databasejournal.java
来自「jsr170接口的java实现。是个apache的开源项目。」· Java 代码 · 共 534 行 · 第 1/2 页
JAVA
534 行
/** * {@inheritDoc} */ protected void doUnlock(boolean successful) { if (!successful) { rollback(con); } try { con.setAutoCommit(true); } catch (SQLException e) { String msg = "Unable to set autocommit to true."; log.warn(msg, e); } } /** * {@inheritDoc} * <p/> * Save away the locked revision inside the newly appended record. */ protected void appending(AppendRecord record) { record.setRevision(lockedRevision); } /** * {@inheritDoc} * <p/> * We have already saved away the revision for this record. */ protected void append(AppendRecord record, InputStream in, int length) throws JournalException { try { try { insertRevisionStmt.clearParameters(); insertRevisionStmt.clearWarnings(); insertRevisionStmt.setLong(1, record.getRevision()); insertRevisionStmt.setString(2, getId()); insertRevisionStmt.setString(3, record.getProducerId()); insertRevisionStmt.setBinaryStream(4, in, length); insertRevisionStmt.execute(); con.commit(); } finally { try { con.setAutoCommit(true); } catch (SQLException e) { String msg = "Unable to set autocommit to true."; log.warn(msg, e); } } } catch (SQLException e) { String msg = "Unable to append revision " + lockedRevision + "."; throw new JournalException(msg, e); } } /** * {@inheritDoc} */ public void close() { try { con.close(); } catch (SQLException e) { String msg = "Error while closing connection: " + e.getMessage(); log.warn(msg); } } /** * Close some input stream. * * @param in input stream, may be <code>null</code>. */ private void close(InputStream in) { try { if (in != null) { in.close(); } } catch (IOException e) { String msg = "Error while closing input stream: " + e.getMessage(); log.warn(msg); } } /** * Close some statement. * * @param stmt statement, may be <code>null</code>. */ private void close(Statement stmt) { try { if (stmt != null) { stmt.close(); } } catch (SQLException e) { String msg = "Error while closing statement: " + e.getMessage(); log.warn(msg); } } /** * Close some resultset. * * @param rs resultset, may be <code>null</code>. */ private void close(ResultSet rs) { try { if (rs != null) { rs.close(); } } catch (SQLException e) { String msg = "Error while closing result set: " + e.getMessage(); log.warn(msg); } } /** * Rollback a connection. * * @param con connection. */ private void rollback(Connection con) { try { con.rollback(); } catch (SQLException e) { String msg = "Error while rolling back connection: " + e.getMessage(); log.warn(msg); } } /** * Checks if the required schema objects exist and creates them if they * don't exist yet. * * @throws Exception if an error occurs */ private void checkSchema() throws Exception { DatabaseMetaData metaData = con.getMetaData(); String tableName = schemaObjectPrefix + "JOURNAL"; if (metaData.storesLowerCaseIdentifiers()) { tableName = tableName.toLowerCase(); } else if (metaData.storesUpperCaseIdentifiers()) { tableName = tableName.toUpperCase(); } ResultSet rs = metaData.getTables(null, null, tableName, null); boolean schemaExists; try { schemaExists = rs.next(); } finally { rs.close(); } if (!schemaExists) { // read ddl from resources InputStream in = DatabaseJournal.class.getResourceAsStream(schema + ".ddl"); if (in == null) { String msg = "No schema-specific DDL found: '" + schema + ".ddl" + "', falling back to '" + DEFAULT_DDL_NAME + "'."; log.info(msg); in = DatabaseJournal.class.getResourceAsStream(DEFAULT_DDL_NAME); if (in == null) { msg = "Unable to load '" + DEFAULT_DDL_NAME + "'."; throw new JournalException(msg); } } BufferedReader reader = new BufferedReader(new InputStreamReader(in)); Statement stmt = con.createStatement(); try { String sql = reader.readLine(); while (sql != null) { // Skip comments and empty lines if (!sql.startsWith("#") && sql.length() > 0) { // replace prefix variable sql = Text.replace(sql, SCHEMA_OBJECT_PREFIX_VARIABLE, schemaObjectPrefix); // execute sql stmt stmt.executeUpdate(sql); } // read next sql stmt sql = reader.readLine(); } } finally { close(in); close(stmt); } } } /** * Builds and prepares the SQL statements. * * @throws SQLException if an error occurs */ private void prepareStatements() throws SQLException { selectRevisionsStmt = con.prepareStatement( "select REVISION_ID, JOURNAL_ID, PRODUCER_ID, REVISION_DATA " + "from " + schemaObjectPrefix + "JOURNAL " + "where REVISION_ID > ?"); updateGlobalStmt = con.prepareStatement( "update " + schemaObjectPrefix + "GLOBAL_REVISION " + "set revision_id = revision_id + 1"); selectGlobalStmt = con.prepareStatement( "select revision_id " + "from " + schemaObjectPrefix + "GLOBAL_REVISION"); insertRevisionStmt = con.prepareStatement( "insert into " + schemaObjectPrefix + "JOURNAL" + "(REVISION_ID, JOURNAL_ID, PRODUCER_ID, REVISION_DATA) " + "values (?,?,?,?)"); } /** * Bean getters */ public String getDriver() { return driver; } public String getUrl() { return url; } public String getSchema() { return schema; } public String getSchemaObjectPrefix() { return schemaObjectPrefix; } public String getUser() { return user; } public String getPassword() { return password; } /** * Bean setters */ public void setDriver(String driver) { this.driver = driver; } public void setUrl(String url) { this.url = url; } public void setSchema(String schema) { this.schema = schema; } public void setSchemaObjectPrefix(String schemaObjectPrefix) { this.schemaObjectPrefix = schemaObjectPrefix.toUpperCase(); } public void setUser(String user) { this.user = user; } public void setPassword(String password) { this.password = password; }}
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?