Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@
import edu.caltech.ipac.table.IpacTableUtil;
import edu.caltech.ipac.table.TableMeta;
import edu.caltech.ipac.table.TableUtil;
import edu.caltech.ipac.table.io.DsvTableIO;
import edu.caltech.ipac.util.FormatUtil;
import edu.caltech.ipac.firefly.server.db.DuckDbReadable;
import edu.caltech.ipac.table.io.IpacTableReader;
import edu.caltech.ipac.table.io.IpacTableWriter;
import edu.caltech.ipac.table.query.DataGroupQuery;
Expand All @@ -34,7 +35,6 @@
import edu.caltech.ipac.util.download.URLDownload;
import edu.caltech.ipac.visualize.plot.CoordinateSys;
import edu.caltech.ipac.visualize.plot.WorldPt;
import org.apache.commons.csv.CSVFormat;

import java.io.BufferedOutputStream;
import java.io.BufferedReader;
Expand Down Expand Up @@ -164,7 +164,7 @@ protected File loadDataFile(TableServerRequest request) throws IOException, Data
// check for errors in returned file
evaluateCVS(csv);

DataGroup dg = DsvTableIO.parse(csv, CSVFormat.DEFAULT.withCommentMarker('#'));
DataGroup dg = csv.length() == 0 ? null : DuckDbReadable.read(FormatUtil.Format.CSV, csv.getAbsolutePath());
if (dg == null) {
_log.info("no data found for search");
return null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -734,35 +734,12 @@ protected Object[] getDdFrom(DataType dt, int colIdx) {
case ROW_NUM -> 1_000_001;
default -> colIdx;
};
return new Object[] {
dt.getKeyName(),
dt.getLabel(),
dt.getTypeDesc(),
dt.getUnits(),
dt.getNullString(),
dt.getFormat(),
dt.getFmtDisp(),
dt.getWidth(),
dt.getVisibility().name(),
dt.isSortable(),
dt.isFilterable(),
dt.isFixed(),
dt.getDesc(),
dt.getEnumVals(),
dt.getID(),
dt.getPrecision(),
dt.getUCD(),
dt.getUType(),
dt.getRef(),
dt.getMaxValue(),
dt.getMinValue(),
Util.serialize(dt.getLinkInfos()), // index(21) is used in HsqlDbAdapter. if it changes, update.
dt.getDataOptions(),
dt.getArraySize(),
dt.getCellRenderer(),
dt.getSortByCols(),
colIdx
};
Object[] row = new Object[EmbeddedDbUtil.DD_COLS.size() + 1];
for (int i = 0; i < EmbeddedDbUtil.DD_COLS.size(); i++) {
row[i] = EmbeddedDbUtil.DD_COLS.get(i).get().apply(dt);
}
row[row.length-1] = colIdx; // order_index is the column's position, not something dt holds
return row;
}

// insert column info into the table
Expand Down Expand Up @@ -964,7 +941,6 @@ Object dbToDD(DataGroup dg, ResultSet rs) {
// if this column is not in DataGroup. no need to update the info
if (dtype != null) {
EmbeddedDbUtil.dbToDataType(dtype, rs);
handleSpecialDTypes(dtype, dg, rs);
}
};
} catch (SQLException e) {
Expand All @@ -973,10 +949,6 @@ Object dbToDD(DataGroup dg, ResultSet rs) {
return 0;
}

void handleSpecialDTypes(DataType dtype, DataGroup dg, ResultSet rs) {
applyIfNotEmpty(deserialize(rs, "links"), v -> dtype.setLinkInfos((List<LinkInfo>) v));
}

public void copyDDFromSource(String tblName, String sourceTbl) {
List<String> cnames = getColumnNamesFromSys(tblName, "'");
String ddSql = "select * from %s_DD".formatted(sourceTbl) +
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@
import edu.caltech.ipac.firefly.server.query.DataAccessException;
import edu.caltech.ipac.firefly.server.util.JsonToDataGroup;
import edu.caltech.ipac.table.DataGroup;
import edu.caltech.ipac.table.io.DsvTableIO;
import edu.caltech.ipac.table.io.FITSTableReader;
import edu.caltech.ipac.table.io.IpacTableReader;
import edu.caltech.ipac.table.io.SpectrumMetaInspector;
Expand Down Expand Up @@ -78,7 +77,7 @@ static FileInfo ingestDuckReadable(FormatUtil.Format format, DbAdapter dbAdapter
} else if (format == PARQUET) {
throw new DataAccessException("Unsupported format (%s), file: %s".formatted(format, source));
} else {
DataGroup table = DsvTableIO.parse(new File(source), format);
DataGroup table = DuckDbReadable.read(format, source); // to avoid using DsvTableIO.parse
return ingestTable(dbAdapter, table, searchForSpectrum);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import edu.caltech.ipac.firefly.server.util.QueryUtil;
import edu.caltech.ipac.firefly.server.util.StopWatch;
import edu.caltech.ipac.table.DataGroup;
import edu.caltech.ipac.table.DataType;
import edu.caltech.ipac.table.io.VoTableReader;
import edu.caltech.ipac.table.io.VoTableWriter;
import edu.caltech.ipac.util.FileUtil;
Expand All @@ -30,6 +31,8 @@
import javax.annotation.Nonnull;

import static edu.caltech.ipac.firefly.core.Util.Opt.ifNotNull;
import static edu.caltech.ipac.firefly.server.db.EmbeddedDbUtil.applyInfoToDataType;
import static edu.caltech.ipac.firefly.server.db.EmbeddedDbUtil.dbToDataGroup;
import static edu.caltech.ipac.table.TableUtil.getAliasName;

/**
Expand Down Expand Up @@ -116,6 +119,49 @@ public DataGroup getInfo(String source) throws DataAccessException {
return table;
}

/** Same as {@link #read(FormatUtil.Format, String, Consumer)}, with no extra meta. */
public static DataGroup read(FormatUtil.Format format, String source) throws DataAccessException {
return read(format, source, null);
}

/**
* Reads the given source with the adapter that handles this format.
* @param format format of the source
* @return the source as a DataGroup, or null if no adapter reads this format
* @see #read(String, Consumer)
*/
public static DataGroup read(FormatUtil.Format format, String source, Consumer<DataGroup> extraMetaSetter) throws DataAccessException {
var adapter = getDetachedAdapter(format);
return adapter == null ? null : adapter.read(source, extraMetaSetter);
}

/** Same as {@link #read(String, Consumer)}, with no extra meta. */
public DataGroup read(String source) throws DataAccessException {
return read(source, null);
}

/**
* Reads the given source into a DataGroup. The data is read into memory; no dbFile needed.
* @param source can be a local file path or a URL
* @param extraMetaSetter additional meta to apply to the returned table
* @return the source as a DataGroup
*/
public DataGroup read(String source, Consumer<DataGroup> extraMetaSetter) throws DataAccessException {
StopWatch.getInstance().start("read: " + source);
DataGroup tableMeta = getTableMeta(source, extraMetaSetter); // the source's schema, with all of its meta applied
String sql = "SELECT * from %s".formatted(sqlReadSource(source));
try {
DataGroup table = getJdbcTmpl().query(sql, rs -> dbToDataGroup(rs, tableMeta)); // adds the rows to it
StopWatch.getInstance().printLog("read: " + source);
return table;
} catch (Exception e) {
LOGGER.warn("read failed with error: " + e.getMessage(),
"sql: " + sql,
"source: " + source);
throw handleSqlExp("Query failed", e);
}
}

/**
* Ingest data directly from a source file. This file can be local or remote.
* @param source can be a local file path or a URL
Expand Down Expand Up @@ -179,20 +225,46 @@ String sqlReadSource(String srcFile) {
return "read_parquet('%s')".formatted(srcFile);
}

/**
* Start from the Parquet schema then apply embedded VOTable metadata on top, ensuring column info matches the actual file.
*/
@Override
protected DataGroup getTableMeta(String source, Consumer<DataGroup> extraMetaSetter) throws DataAccessException {
DataGroup tableMeta = super.getTableMeta(source, null); // the schema, straight from the file
ifNotNull(readVoTableMeta(source)).apply(voMeta -> applyVoMeta(tableMeta, voMeta));
if (extraMetaSetter != null) extraMetaSetter.accept(tableMeta);
return tableMeta;
}

/**
* Applies VOTable metadata onto the given table
*/
private static void applyVoMeta(DataGroup table, DataGroup voMeta) {
for (DataType col : table.getDataDefinitions()) {
DataType info = voMeta.getDataDefintion(col.getKeyName()); // exact match first;
if (info == null) info = voMeta.getDataDefintion(col.getKeyName(), true); // then, ignore case
applyInfoToDataType(col, info);
}
table.setTitle(voMeta.getTitle());
table.setTableMeta(voMeta.getTableMeta());
table.setGroupInfos(voMeta.getGroupInfos());
table.setLinkInfos(voMeta.getLinkInfos());
table.setParamInfos(voMeta.getParamInfos());
table.setResourceInfos(voMeta.getResourceInfos());
}

/**
* @return the VOTable stored in the Parquet metadata, or null if unavailable or unreadable
*/
private DataGroup readVoTableMeta(String source) {
var jdbc = JdbcFactory.getTemplate(getDbInstance());
try {
var votable = jdbc.queryForObject(
"SELECT decode(value) FROM parquet_kv_metadata('%s') where key = 'IVOA.VOTable-Parquet.content'".formatted(source),
String.class);
if (votable != null) {
DataGroup tableMeta = VoTableReader.voToDataGroups(new ByteArrayInputStream(votable.getBytes()), false)[0];
if (tableMeta != null && extraMetaSetter != null) extraMetaSetter.accept(tableMeta);
return tableMeta;
}
} catch (Exception ignored) {} // ignored if it can't read
return super.getTableMeta(source, extraMetaSetter);
return votable == null ? null :
VoTableReader.voToDataGroups(new ByteArrayInputStream(votable.getBytes()), false)[0];
} catch (Exception ignored) { return null; } // ignored if it can't read
}

public void export(TableServerRequest treq, OutputStream out) throws DataAccessException {
Expand Down
Loading
Loading