Skip to content
Open
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
11 changes: 11 additions & 0 deletions iotdb-client/client-cpp/src/include/Session.h
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,8 @@ template <typename T, typename Target> void safe_cast(const T& value, Target& ta
*
*/
class Tablet {
friend class SessionUtils;

private:
static const int DEFAULT_ROW_SIZE = 1024;

Expand Down Expand Up @@ -431,6 +433,15 @@ class SessionUtils {
static std::string getValue(const Tablet& tablet);

static bool isTabletContainsSingleDevice(Tablet tablet);

/**
* Drop entirely-null FIELD columns within [0, rowSize). TAG/ATTRIBUTE are kept.
* Does not mutate {@code tablet}.
*
* @return non-owning pointer to {@code tablet} if nothing to drop; a new tablet if
* filtered; nullptr if every FIELD column is null (caller should skip insert)
*/
static std::shared_ptr<const Tablet> filterNullColumns(const Tablet& tablet);
};

class TemplateNode {
Expand Down
3 changes: 2 additions & 1 deletion iotdb-client/client-cpp/src/rpc/SessionImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -145,7 +145,8 @@ class Session::Impl {
void handleRedirection(const std::string& deviceId, TEndPoint endPoint);
void handleRedirection(const std::shared_ptr<storage::IDeviceID>& deviceId, TEndPoint endPoint);

static void buildInsertTabletReq(TSInsertTabletReq& request, Tablet& tablet, bool sorted);
// Returns false when all FIELD columns are null and the insert should be skipped.
static bool buildInsertTabletReq(TSInsertTabletReq& request, Tablet& tablet, bool sorted);
Comment on lines +148 to +149

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

For the table model, time-only insertion is allowed.

void insertTablet(TSInsertTabletReq request);
void insertRelationalTabletOnce(
const std::unordered_map<std::shared_ptr<SessionConnection>, Tablet>& relationalTabletGroup,
Expand Down
167 changes: 134 additions & 33 deletions iotdb-client/client-cpp/src/session/Session.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -413,6 +413,83 @@ bool SessionUtils::isTabletContainsSingleDevice(Tablet tablet) {
return true;
}

static bool isColumnAllNull(const BitMap& bitMap, size_t rowSize) {
if (rowSize == 0) {
return false;
}
if (bitMap.getSize() == rowSize && bitMap.isAllMarked()) {
return true;
}
for (size_t row = 0; row < rowSize; row++) {
if (!bitMap.isMarked(row)) {
return false;
}
}
return true;
}
Comment on lines +416 to +429

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why is bitMap.isAllMarked not enough?


std::shared_ptr<const Tablet> SessionUtils::filterNullColumns(const Tablet& tablet) {
const size_t columnCount = tablet.schemas.size();
if (columnCount == 0 || tablet.bitMaps.size() < columnCount) {
return std::shared_ptr<const Tablet>(&tablet, [](const Tablet*) {});
}

std::vector<size_t> keptIndices;
keptIndices.reserve(columnCount);
size_t originalFieldCount = 0;
size_t keptFieldCount = 0;

for (size_t i = 0; i < columnCount; i++) {
ColumnCategory category =
i < tablet.columnTypes.size() ? tablet.columnTypes[i] : ColumnCategory::FIELD;
bool isField = category == ColumnCategory::FIELD;
if (isField) {
originalFieldCount++;
}
bool drop = isField && isColumnAllNull(tablet.bitMaps[i], tablet.rowSize);
if (drop) {
continue;
}
keptIndices.push_back(i);
if (isField) {
keptFieldCount++;
}
}

if (keptIndices.size() == columnCount) {
return std::shared_ptr<const Tablet>(&tablet, [](const Tablet*) {});
}
if (originalFieldCount > 0 && keptFieldCount == 0) {
return nullptr;
}
if (keptIndices.empty()) {
return nullptr;
}

std::vector<std::pair<std::string, TSDataType::TSDataType>> keptSchemas;
std::vector<ColumnCategory> keptColumnTypes;
keptSchemas.reserve(keptIndices.size());
keptColumnTypes.reserve(keptIndices.size());
for (size_t idx : keptIndices) {
keptSchemas.push_back(tablet.schemas[idx]);
keptColumnTypes.push_back(idx < tablet.columnTypes.size() ? tablet.columnTypes[idx]
: ColumnCategory::FIELD);
}

auto filteredOut = std::make_shared<Tablet>(tablet.deviceId, keptSchemas, keptColumnTypes,
tablet.maxRowNumber, tablet.isAligned);
filteredOut->deleteColumns();
Comment on lines +479 to +481

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not create-and-delete or copy. May add a constructor for this.

filteredOut->timestamps = tablet.timestamps;
filteredOut->rowSize = tablet.rowSize;
for (size_t ni = 0; ni < keptIndices.size(); ni++) {
size_t oi = keptIndices[ni];
Tablet::deepCopyTabletColValue(&tablet.values[oi], &filteredOut->values[ni],
keptSchemas[ni].second, static_cast<int>(tablet.maxRowNumber));
filteredOut->bitMaps[ni] = tablet.bitMaps[oi];
}
return filteredOut;
}

string MeasurementNode::serialize() const {
MyStringBuffer buffer;
buffer.putString(getName());
Expand Down Expand Up @@ -1002,27 +1079,35 @@ void Session::Impl::insertTabletsWithLeaderCache(unordered_map<string, Tablet*>&
}
auto deviceId = item.first;
auto tablet = item.second;
std::shared_ptr<const Tablet> toEncode = SessionUtils::filterNullColumns(*tablet);
if (!toEncode) {
continue;
}
auto connection = getSessionConnection(deviceId);
auto it = tabletsGroup.find(connection);
if (it == tabletsGroup.end()) {
TSInsertTabletsReq request;
tabletsGroup[connection] = request;
}
TSInsertTabletsReq& existingReq = tabletsGroup[connection];
existingReq.prefixPaths.emplace_back(tablet->deviceId);
existingReq.timestampsList.emplace_back(move(SessionUtils::getTime(*tablet)));
existingReq.valuesList.emplace_back(move(SessionUtils::getValue(*tablet)));
existingReq.sizeList.emplace_back(tablet->rowSize);
existingReq.prefixPaths.emplace_back(toEncode->deviceId);
existingReq.timestampsList.emplace_back(move(SessionUtils::getTime(*toEncode)));
existingReq.valuesList.emplace_back(move(SessionUtils::getValue(*toEncode)));
existingReq.sizeList.emplace_back(toEncode->rowSize);
vector<int> dataTypes;
vector<string> measurements;
for (pair<string, TSDataType::TSDataType> schema : tablet->schemas) {
for (pair<string, TSDataType::TSDataType> schema : toEncode->schemas) {
measurements.push_back(schema.first);
dataTypes.push_back(schema.second);
}
existingReq.measurementsList.emplace_back(measurements);
existingReq.typesList.emplace_back(dataTypes);
}

if (tabletsGroup.empty()) {
return;
}

std::function<void(std::shared_ptr<SessionConnection>, const TSInsertTabletsReq&)> consumer =
[](const std::shared_ptr<SessionConnection>& c, const TSInsertTabletsReq& r) {
c->insertTablets(r);
Expand Down Expand Up @@ -1444,27 +1529,41 @@ void Session::insertTablet(Tablet& tablet) {
}
}

void Session::Impl::buildInsertTabletReq(TSInsertTabletReq& request, Tablet& tablet, bool sorted) {
bool Session::Impl::buildInsertTabletReq(TSInsertTabletReq& request, Tablet& tablet, bool sorted) {
if ((!sorted) && !checkSorted(tablet)) {
sortTablet(tablet);
}

request.__set_prefixPath(tablet.deviceId);
std::shared_ptr<const Tablet> toEncode = SessionUtils::filterNullColumns(tablet);
if (!toEncode) {
return false;
}

request.__set_prefixPath(toEncode->deviceId);

std::vector<std::string> reqMeasurements;
reqMeasurements.reserve(tablet.schemas.size());
reqMeasurements.reserve(toEncode->schemas.size());
std::vector<int32_t> types;
types.reserve(tablet.schemas.size());
for (pair<string, TSDataType::TSDataType> schema : tablet.schemas) {
types.reserve(toEncode->schemas.size());
for (pair<string, TSDataType::TSDataType> schema : toEncode->schemas) {
reqMeasurements.push_back(schema.first);
types.push_back(schema.second);
}
request.__set_measurements(reqMeasurements);
request.__set_types(types);
request.__set_values(SessionUtils::getValue(tablet));
request.__set_timestamps(SessionUtils::getTime(tablet));
request.__set_size(tablet.rowSize);
request.__set_isAligned(tablet.isAligned);
request.__set_values(SessionUtils::getValue(*toEncode));
request.__set_timestamps(SessionUtils::getTime(*toEncode));
request.__set_size(toEncode->rowSize);
request.__set_isAligned(toEncode->isAligned);
if (!toEncode->columnTypes.empty()) {
std::vector<int8_t> columnCategories;
columnCategories.reserve(toEncode->columnTypes.size());
for (auto& category : toEncode->columnTypes) {
columnCategories.push_back(static_cast<int8_t>(category));
}
request.__set_columnCategories(columnCategories);
}
return true;
}

void Session::Impl::insertTablet(TSInsertTabletReq request) {
Expand All @@ -1488,7 +1587,9 @@ void Session::Impl::insertTablet(TSInsertTabletReq request) {

void Session::insertTablet(Tablet& tablet, bool sorted) {
TSInsertTabletReq request;
impl_->buildInsertTabletReq(request, tablet, sorted);
if (!impl_->buildInsertTabletReq(request, tablet, sorted)) {
return;
}
impl_->insertTablet(request);
}

Expand Down Expand Up @@ -1574,13 +1675,10 @@ void Session::Impl::insertRelationalTabletOnce(
auto connection = iter->first;
auto tablet = iter->second;
TSInsertTabletReq request;
buildInsertTabletReq(request, tablet, sorted);
request.__set_writeToTable(true);
std::vector<int8_t> columnCategories;
for (auto& category : tablet.columnTypes) {
columnCategories.push_back(static_cast<int8_t>(category));
if (!buildInsertTabletReq(request, tablet, sorted)) {
return;
}
request.__set_columnCategories(columnCategories);
request.__set_writeToTable(true);
try {
TSStatus respStatus;
connection->getSessionClient()->insertTablet(respStatus, request);
Expand Down Expand Up @@ -1629,14 +1727,10 @@ void Session::Impl::insertRelationalTabletByGroup(
futures.emplace_back(
std::async(std::launch::async, [this, connection, tablet, sorted]() mutable {
TSInsertTabletReq request;
buildInsertTabletReq(request, tablet, sorted);
request.__set_writeToTable(true);

std::vector<int8_t> columnCategories;
for (auto& category : tablet.columnTypes) {
columnCategories.push_back(static_cast<int8_t>(category));
if (!buildInsertTabletReq(request, tablet, sorted)) {
return;
}
request.__set_columnCategories(columnCategories);
request.__set_writeToTable(true);

try {
TSStatus respStatus;
Expand Down Expand Up @@ -1702,18 +1796,25 @@ void Session::insertTablets(unordered_map<string, Tablet*>& tablets, bool sorted
if (!impl_->checkSorted(*(item.second))) {
impl_->sortTablet(*(item.second));
}
request.prefixPaths.push_back(item.second->deviceId);
std::shared_ptr<const Tablet> toEncode = SessionUtils::filterNullColumns(*(item.second));
if (!toEncode) {
continue;
}
request.prefixPaths.push_back(toEncode->deviceId);
vector<string> measurements;
vector<int> dataTypes;
for (pair<string, TSDataType::TSDataType> schema : item.second->schemas) {
for (pair<string, TSDataType::TSDataType> schema : toEncode->schemas) {
measurements.push_back(schema.first);
dataTypes.push_back(schema.second);
}
request.measurementsList.push_back(measurements);
request.typesList.push_back(dataTypes);
request.timestampsList.push_back(move(SessionUtils::getTime(*(item.second))));
request.valuesList.push_back(move(SessionUtils::getValue(*(item.second))));
request.sizeList.push_back(item.second->rowSize);
request.timestampsList.push_back(move(SessionUtils::getTime(*toEncode)));
request.valuesList.push_back(move(SessionUtils::getValue(*toEncode)));
request.sizeList.push_back(toEncode->rowSize);
}
if (request.prefixPaths.empty()) {
return;
}
request.__set_isAligned(isAligned);
try {
Expand Down
7 changes: 6 additions & 1 deletion iotdb-client/client-cpp/test/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -42,12 +42,14 @@ set(_test_targets
session_tests
session_relational_tests
session_c_tests
session_c_relational_tests)
session_c_relational_tests
session_utils_tests)

add_executable(session_tests main.cpp cpp/sessionIT.cpp)
add_executable(session_relational_tests main_Relational.cpp cpp/sessionRelationalIT.cpp)
add_executable(session_c_tests main_c.cpp cpp/sessionCIT.cpp)
add_executable(session_c_relational_tests main_c_Relational.cpp cpp/sessionCRelationalIT.cpp)
add_executable(session_utils_tests main_utils.cpp cpp/sessionUtilsTest.cpp)

foreach(_t IN LISTS _test_targets)
target_include_directories(${_t} PRIVATE
Expand Down Expand Up @@ -86,6 +88,7 @@ if(MSVC)
add_test(NAME sessionRelationalIT CONFIGURATIONS Release COMMAND session_relational_tests)
add_test(NAME sessionCIT CONFIGURATIONS Release COMMAND session_c_tests)
add_test(NAME sessionCRelationalIT CONFIGURATIONS Release COMMAND session_c_relational_tests)
add_test(NAME sessionUtilsTest CONFIGURATIONS Release COMMAND session_utils_tests)
foreach(_t IN LISTS _test_targets)
add_custom_command(TARGET ${_t} POST_BUILD
COMMAND ${CMAKE_COMMAND} -E copy_if_different
Expand All @@ -96,9 +99,11 @@ else()
add_test(NAME sessionRelationalIT COMMAND session_relational_tests)
add_test(NAME sessionCIT COMMAND session_c_tests)
add_test(NAME sessionCRelationalIT COMMAND session_c_relational_tests)
add_test(NAME sessionUtilsTest COMMAND session_utils_tests)
endif()

# Run sequentially: parallel ctest overloads the single local IoTDB instance.
# sessionUtilsTest is a pure unit test and can run anytime.
set_tests_properties(
sessionIT sessionRelationalIT sessionCIT sessionCRelationalIT
PROPERTIES RUN_SERIAL TRUE)
Loading
Loading