[MM-14253] Adds channels and users to the bulk index process (#10434)
* [MM-14253] Adds channels and users to the bulk index process * Add support for PostgreSQL and sort the user query result * Add tests for user and channel batch queries * Fix test times
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
5a9d95d9c7
Коммит
372ef87f76
@@ -2583,3 +2583,32 @@ func (s SqlChannelStore) GetAllDirectChannelsForExportAfter(limit int, afterId s
|
||||
result.Data = directChannelsForExport
|
||||
})
|
||||
}
|
||||
|
||||
func (s SqlChannelStore) GetChannelsBatchForIndexing(startTime, endTime int64, limit int) store.StoreChannel {
|
||||
return store.Do(func(result *store.StoreResult) {
|
||||
var channels []*model.Channel
|
||||
_, err1 := s.GetSearchReplica().Select(&channels,
|
||||
`SELECT
|
||||
*
|
||||
FROM
|
||||
Channels
|
||||
WHERE
|
||||
Type = 'O'
|
||||
AND
|
||||
CreateAt >= :StartTime
|
||||
AND
|
||||
CreateAt < :EndTime
|
||||
ORDER BY
|
||||
CreateAt
|
||||
LIMIT
|
||||
:NumChannels`,
|
||||
map[string]interface{}{"StartTime": startTime, "EndTime": endTime, "NumChannels": limit})
|
||||
|
||||
if err1 != nil {
|
||||
result.Err = model.NewAppError("SqlChannelStore.GetChannelsBatchForIndexing", "store.sql_channel.get_channels_batch_for_indexing.get.app_error", nil, err1.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
result.Data = channels
|
||||
})
|
||||
}
|
||||
|
||||
@@ -269,14 +269,14 @@ func (s *SqlPostStore) GetFlaggedPostsForChannel(userId, channelId string, offse
|
||||
|
||||
var posts []*model.Post
|
||||
query := `
|
||||
SELECT
|
||||
*
|
||||
FROM Posts
|
||||
WHERE
|
||||
Id IN (SELECT Name FROM Preferences WHERE UserId = :UserId AND Category = :Category)
|
||||
SELECT
|
||||
*
|
||||
FROM Posts
|
||||
WHERE
|
||||
Id IN (SELECT Name FROM Preferences WHERE UserId = :UserId AND Category = :Category)
|
||||
AND ChannelId = :ChannelId
|
||||
AND DeleteAt = 0
|
||||
ORDER BY CreateAt DESC
|
||||
AND DeleteAt = 0
|
||||
ORDER BY CreateAt DESC
|
||||
LIMIT :Limit OFFSET :Offset`
|
||||
|
||||
if _, err := s.GetReplica().Select(&posts, query, map[string]interface{}{"UserId": userId, "Category": model.PREFERENCE_CATEGORY_FLAGGED_POST, "ChannelId": channelId, "Offset": offset, "Limit": limit}); err != nil {
|
||||
@@ -835,7 +835,7 @@ func (s *SqlPostStore) Search(teamId string, userId string, params *model.Search
|
||||
` + userIdPart + `
|
||||
` + deletedQueryPart + `
|
||||
CHANNEL_FILTER)
|
||||
CREATEDATE_CLAUSE
|
||||
CREATEDATE_CLAUSE
|
||||
SEARCH_CLAUSE
|
||||
ORDER BY CreateAt DESC
|
||||
LIMIT 100`
|
||||
@@ -1264,11 +1264,11 @@ func (s *SqlPostStore) determineMaxPostSize() int {
|
||||
// The Post.Message column in MySQL has historically been TEXT, with a maximum
|
||||
// limit of 65535.
|
||||
if err := s.GetReplica().SelectOne(&maxPostSizeBytes, `
|
||||
SELECT
|
||||
SELECT
|
||||
COALESCE(CHARACTER_MAXIMUM_LENGTH, 0)
|
||||
FROM
|
||||
FROM
|
||||
INFORMATION_SCHEMA.COLUMNS
|
||||
WHERE
|
||||
WHERE
|
||||
table_schema = DATABASE()
|
||||
AND table_name = 'Posts'
|
||||
AND column_name = 'Message'
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
sq "github.com/Masterminds/squirrel"
|
||||
@@ -1363,15 +1364,15 @@ func (us SqlUserStore) GetEtagForProfilesNotInTeam(teamId string) store.StoreCha
|
||||
|
||||
var querystr string
|
||||
querystr = `
|
||||
SELECT
|
||||
SELECT
|
||||
CONCAT(MAX(UpdateAt), '.', COUNT(Id)) as etag
|
||||
FROM
|
||||
FROM
|
||||
Users as u
|
||||
LEFT JOIN TeamMembers tm
|
||||
ON tm.UserId = u.Id
|
||||
AND tm.TeamId = :TeamId
|
||||
LEFT JOIN TeamMembers tm
|
||||
ON tm.UserId = u.Id
|
||||
AND tm.TeamId = :TeamId
|
||||
AND tm.DeleteAt = 0
|
||||
WHERE
|
||||
WHERE
|
||||
tm.UserId IS NULL
|
||||
`
|
||||
etag, err := us.GetReplica().SelectStr(querystr, map[string]interface{}{"TeamId": teamId})
|
||||
@@ -1449,3 +1450,89 @@ func (us SqlUserStore) InferSystemInstallDate() store.StoreChannel {
|
||||
result.Data = createAt
|
||||
})
|
||||
}
|
||||
|
||||
func (us SqlUserStore) GetUsersBatchForIndexing(startTime, endTime int64, limit int) store.StoreChannel {
|
||||
return store.Do(func(result *store.StoreResult) {
|
||||
var users []*model.User
|
||||
usersQuery, args, _ := us.usersQuery.
|
||||
Where(sq.GtOrEq{"u.CreateAt": startTime}).
|
||||
Where(sq.Lt{"u.CreateAt": endTime}).
|
||||
OrderBy("u.CreateAt").
|
||||
Limit(uint64(limit)).
|
||||
ToSql()
|
||||
_, err1 := us.GetSearchReplica().Select(&users, usersQuery, args...)
|
||||
|
||||
if err1 != nil {
|
||||
result.Err = model.NewAppError("SqlUserStore.GetUsersBatchForIndexing", "store.sql_user.get_users_batch_for_indexing.get_users.app_error", nil, err1.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
userIds := []string{}
|
||||
for _, user := range users {
|
||||
userIds = append(userIds, user.Id)
|
||||
}
|
||||
|
||||
var channelMembers []*model.ChannelMember
|
||||
channelMembersQuery, args, _ := us.getQueryBuilder().
|
||||
Select("cm.*").
|
||||
From("ChannelMembers cm").
|
||||
Join("Channels c ON cm.ChannelId = c.Id").
|
||||
Where(sq.Eq{"c.Type": "O", "cm.UserId": userIds}).
|
||||
ToSql()
|
||||
_, err2 := us.GetSearchReplica().Select(&channelMembers, channelMembersQuery, args...)
|
||||
|
||||
if err2 != nil {
|
||||
result.Err = model.NewAppError("SqlUserStore.GetUsersBatchForIndexing", "store.sql_user.get_users_batch_for_indexing.get_channel_members.app_error", nil, err2.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
var teamMembers []*model.TeamMember
|
||||
teamMembersQuery, args, _ := us.getQueryBuilder().
|
||||
Select("*").
|
||||
From("TeamMembers").
|
||||
Where(sq.Eq{"UserId": userIds, "DeleteAt": 0}).
|
||||
ToSql()
|
||||
_, err3 := us.GetSearchReplica().Select(&teamMembers, teamMembersQuery, args...)
|
||||
|
||||
if err3 != nil {
|
||||
result.Err = model.NewAppError("SqlUserStore.GetUsersBatchForIndexing", "store.sql_user.get_users_batch_for_indexing.get_team_members.app_error", nil, err3.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
userMap := map[string]*model.UserForIndexing{}
|
||||
for _, user := range users {
|
||||
userMap[user.Id] = &model.UserForIndexing{
|
||||
Id: user.Id,
|
||||
Username: user.Username,
|
||||
Nickname: user.Nickname,
|
||||
FirstName: user.FirstName,
|
||||
LastName: user.LastName,
|
||||
CreateAt: user.CreateAt,
|
||||
DeleteAt: user.DeleteAt,
|
||||
TeamsIds: []string{},
|
||||
ChannelsIds: []string{},
|
||||
}
|
||||
}
|
||||
|
||||
for _, c := range channelMembers {
|
||||
if userMap[c.UserId] != nil {
|
||||
userMap[c.UserId].ChannelsIds = append(userMap[c.UserId].ChannelsIds, c.ChannelId)
|
||||
}
|
||||
}
|
||||
for _, t := range teamMembers {
|
||||
if userMap[t.UserId] != nil {
|
||||
userMap[t.UserId].TeamsIds = append(userMap[t.UserId].TeamsIds, t.TeamId)
|
||||
}
|
||||
}
|
||||
|
||||
usersForIndexing := []*model.UserForIndexing{}
|
||||
for _, user := range userMap {
|
||||
usersForIndexing = append(usersForIndexing, user)
|
||||
}
|
||||
sort.Slice(usersForIndexing, func(i, j int) bool {
|
||||
return usersForIndexing[i].CreateAt < usersForIndexing[j].CreateAt
|
||||
})
|
||||
|
||||
result.Data = usersForIndexing
|
||||
})
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user